use percent_encoding::{percent_decode, percent_decode_str};
use url::Url;
use std::{
borrow::Cow, collections::HashMap, hash::Hash, net::SocketAddr, path::Path, str::FromStr,
time::Duration,
};
use crate::{consts::CapabilityFlags, LocalInfileHandler, UrlError};
pub const DEFAULT_STMT_CACHE_SIZE: usize = 32;
#[derive(Debug, Clone, Eq, PartialEq, Hash, Default)]
pub struct SslOpts {
pkcs12_path: Option<Cow<'static, Path>>,
password: Option<Cow<'static, str>>,
root_cert_path: Option<Cow<'static, Path>>,
skip_domain_validation: bool,
accept_invalid_certs: bool,
}
impl SslOpts {
pub fn with_pkcs12_path<T: Into<Cow<'static, Path>>>(mut self, pkcs12_path: Option<T>) -> Self {
self.pkcs12_path = pkcs12_path.map(Into::into);
self
}
pub fn with_password<T: Into<Cow<'static, str>>>(mut self, password: Option<T>) -> Self {
self.password = password.map(Into::into);
self
}
pub fn with_root_cert_path<T: Into<Cow<'static, Path>>>(
mut self,
root_cert_path: Option<T>,
) -> Self {
self.root_cert_path = root_cert_path.map(Into::into);
self
}
pub fn with_danger_skip_domain_validation(mut self, value: bool) -> Self {
self.skip_domain_validation = value;
self
}
pub fn with_danger_accept_invalid_certs(mut self, value: bool) -> Self {
self.accept_invalid_certs = value;
self
}
pub fn pkcs12_path(&self) -> Option<&Path> {
self.pkcs12_path.as_ref().map(|x| x.as_ref())
}
pub fn password(&self) -> Option<&str> {
self.password.as_ref().map(AsRef::as_ref)
}
pub fn root_cert_path(&self) -> Option<&Path> {
self.root_cert_path.as_ref().map(AsRef::as_ref)
}
pub fn skip_domain_validation(&self) -> bool {
self.skip_domain_validation
}
pub fn accept_invalid_certs(&self) -> bool {
self.accept_invalid_certs
}
}
#[derive(Debug, Clone, Eq, PartialEq)]
pub(crate) struct InnerOpts {
ip_or_hostname: url::Host,
tcp_port: u16,
socket: Option<String>,
user: Option<String>,
pass: Option<String>,
db_name: Option<String>,
read_timeout: Option<Duration>,
write_timeout: Option<Duration>,
prefer_socket: bool,
tcp_nodelay: bool,
tcp_keepalive_time: Option<u32>,
init: Vec<String>,
ssl_opts: Option<SslOpts>,
local_infile_handler: Option<LocalInfileHandler>,
tcp_connect_timeout: Option<Duration>,
bind_address: Option<SocketAddr>,
stmt_cache_size: usize,
compress: Option<crate::Compression>,
additional_capabilities: CapabilityFlags,
connect_attrs: HashMap<String, String>,
#[cfg(test)]
pub injected_socket: Option<String>,
}
impl Default for InnerOpts {
fn default() -> Self {
InnerOpts {
ip_or_hostname: url::Host::Domain(String::from("localhost")),
tcp_port: 3306,
socket: None,
user: None,
pass: None,
db_name: None,
read_timeout: None,
write_timeout: None,
prefer_socket: true,
init: vec![],
ssl_opts: None,
tcp_keepalive_time: None,
tcp_nodelay: true,
local_infile_handler: None,
tcp_connect_timeout: None,
bind_address: None,
stmt_cache_size: DEFAULT_STMT_CACHE_SIZE,
compress: None,
additional_capabilities: CapabilityFlags::empty(),
connect_attrs: HashMap::new(),
#[cfg(test)]
injected_socket: None,
}
}
}
#[derive(Clone, Eq, PartialEq, Debug, Default)]
pub struct Opts(pub(crate) Box<InnerOpts>);
impl Opts {
#[doc(hidden)]
pub fn addr_is_loopback(&self) -> bool {
match self.0.ip_or_hostname {
url::Host::Domain(ref name) => name == "localhost",
url::Host::Ipv4(ref addr) => addr.is_loopback(),
url::Host::Ipv6(ref addr) => addr.is_loopback(),
}
}
pub fn from_url(url: &str) -> Result<Opts, UrlError> {
from_url(url)
}
pub(crate) fn get_host(&self) -> url::Host {
self.0.ip_or_hostname.clone()
}
pub fn get_ip_or_hostname(&self) -> Cow<str> {
self.0.ip_or_hostname.to_string().into()
}
pub fn get_tcp_port(&self) -> u16 {
self.0.tcp_port
}
pub fn get_socket(&self) -> Option<&str> {
self.0.socket.as_ref().map(|x| &**x)
}
pub fn get_user(&self) -> Option<&str> {
self.0.user.as_ref().map(|x| &**x)
}
pub fn get_pass(&self) -> Option<&str> {
self.0.pass.as_ref().map(|x| &**x)
}
pub fn get_db_name(&self) -> Option<&str> {
self.0.db_name.as_ref().map(|x| &**x)
}
pub fn get_read_timeout(&self) -> Option<&Duration> {
self.0.read_timeout.as_ref()
}
pub fn get_write_timeout(&self) -> Option<&Duration> {
self.0.write_timeout.as_ref()
}
pub fn get_prefer_socket(&self) -> bool {
self.0.prefer_socket
}
pub fn get_init(&self) -> Vec<String> {
self.0.init.clone()
}
pub fn get_ssl_opts(&self) -> Option<&SslOpts> {
self.0.ssl_opts.as_ref()
}
fn set_prefer_socket(&mut self, val: bool) {
self.0.prefer_socket = val;
}
pub fn get_tcp_nodelay(&self) -> bool {
self.0.tcp_nodelay
}
pub fn get_tcp_keepalive_time_ms(&self) -> Option<u32> {
self.0.tcp_keepalive_time
}
pub fn get_local_infile_handler(&self) -> Option<&LocalInfileHandler> {
self.0.local_infile_handler.as_ref()
}
pub fn get_tcp_connect_timeout(&self) -> Option<Duration> {
self.0.tcp_connect_timeout
}
pub fn bind_address(&self) -> Option<&SocketAddr> {
self.0.bind_address.as_ref()
}
pub fn get_stmt_cache_size(&self) -> usize {
self.0.stmt_cache_size
}
pub fn get_compress(&self) -> Option<crate::Compression> {
self.0.compress
}
pub fn get_additional_capabilities(&self) -> CapabilityFlags {
self.0.additional_capabilities
}
pub fn get_connect_attrs(&self) -> &HashMap<String, String> {
&self.0.connect_attrs
}
}
#[derive(Debug)]
pub struct OptsBuilder {
opts: Opts,
}
impl OptsBuilder {
pub fn new() -> Self {
OptsBuilder::default()
}
pub fn from_opts<T: Into<Opts>>(opts: T) -> Self {
OptsBuilder { opts: opts.into() }
}
pub fn ip_or_hostname<T: Into<String>>(mut self, ip_or_hostname: Option<T>) -> Self {
let new = ip_or_hostname.map(Into::into).unwrap_or("127.0.0.1".into());
self.opts.0.ip_or_hostname = url::Host::parse(&new)
.map(|host| host.to_owned())
.unwrap_or_else(|_| url::Host::Domain(new.to_owned()));
self
}
pub fn tcp_port(mut self, tcp_port: u16) -> Self {
self.opts.0.tcp_port = tcp_port;
self
}
pub fn socket<T: Into<String>>(mut self, socket: Option<T>) -> Self {
self.opts.0.socket = socket.map(Into::into);
self
}
pub fn user<T: Into<String>>(mut self, user: Option<T>) -> Self {
self.opts.0.user = user.map(Into::into);
self
}
pub fn pass<T: Into<String>>(mut self, pass: Option<T>) -> Self {
self.opts.0.pass = pass.map(Into::into);
self
}
pub fn db_name<T: Into<String>>(mut self, db_name: Option<T>) -> Self {
self.opts.0.db_name = db_name.map(Into::into);
self
}
pub fn read_timeout(mut self, read_timeout: Option<Duration>) -> Self {
self.opts.0.read_timeout = read_timeout;
self
}
pub fn write_timeout(mut self, write_timeout: Option<Duration>) -> Self {
self.opts.0.write_timeout = write_timeout;
self
}
pub fn tcp_keepalive_time_ms(mut self, tcp_keepalive_time_ms: Option<u32>) -> Self {
self.opts.0.tcp_keepalive_time = tcp_keepalive_time_ms;
self
}
pub fn tcp_nodelay(mut self, nodelay: bool) -> Self {
self.opts.0.tcp_nodelay = nodelay;
self
}
pub fn prefer_socket(mut self, prefer_socket: bool) -> Self {
self.opts.0.prefer_socket = prefer_socket;
self
}
pub fn init<T: Into<String>>(mut self, init: Vec<T>) -> Self {
self.opts.0.init = init.into_iter().map(Into::into).collect();
self
}
pub fn ssl_opts<T: Into<Option<SslOpts>>>(mut self, ssl_opts: T) -> Self {
self.opts.0.ssl_opts = ssl_opts.into();
self
}
pub fn local_infile_handler(mut self, handler: Option<LocalInfileHandler>) -> Self {
self.opts.0.local_infile_handler = handler;
self
}
pub fn tcp_connect_timeout(mut self, timeout: Option<Duration>) -> Self {
self.opts.0.tcp_connect_timeout = timeout;
self
}
pub fn bind_address<T>(mut self, bind_address: Option<T>) -> Self
where
T: Into<SocketAddr>,
{
self.opts.0.bind_address = bind_address.map(Into::into);
self
}
pub fn stmt_cache_size<T>(mut self, cache_size: T) -> Self
where
T: Into<Option<usize>>,
{
self.opts.0.stmt_cache_size = cache_size.into().unwrap_or(128);
self
}
pub fn compress(mut self, compress: Option<crate::Compression>) -> Self {
self.opts.0.compress = compress;
self
}
pub fn additional_capabilities(mut self, additional_capabilities: CapabilityFlags) -> Self {
let forbidden_flags: CapabilityFlags = CapabilityFlags::CLIENT_PROTOCOL_41
| CapabilityFlags::CLIENT_SSL
| CapabilityFlags::CLIENT_COMPRESS
| CapabilityFlags::CLIENT_SECURE_CONNECTION
| CapabilityFlags::CLIENT_LONG_PASSWORD
| CapabilityFlags::CLIENT_TRANSACTIONS
| CapabilityFlags::CLIENT_LOCAL_FILES
| CapabilityFlags::CLIENT_MULTI_STATEMENTS
| CapabilityFlags::CLIENT_MULTI_RESULTS
| CapabilityFlags::CLIENT_PS_MULTI_RESULTS;
self.opts.0.additional_capabilities = additional_capabilities & !forbidden_flags;
self
}
pub fn connect_attrs<T1: Into<String> + Eq + Hash, T2: Into<String>>(
mut self,
connect_attrs: HashMap<T1, T2>,
) -> Self {
self.opts.0.connect_attrs = HashMap::with_capacity(connect_attrs.len());
for (name, value) in connect_attrs {
let name = name.into();
if !name.starts_with('_') {
self.opts.0.connect_attrs.insert(name, value.into());
}
}
self
}
}
impl From<OptsBuilder> for Opts {
fn from(builder: OptsBuilder) -> Opts {
builder.opts
}
}
impl Default for OptsBuilder {
fn default() -> OptsBuilder {
OptsBuilder {
opts: Opts::default(),
}
}
}
fn get_opts_user_from_url(url: &Url) -> Option<String> {
let user = url.username();
if user != "" {
Some(
percent_decode(user.as_ref())
.decode_utf8_lossy()
.into_owned(),
)
} else {
None
}
}
fn get_opts_pass_from_url(url: &Url) -> Option<String> {
if let Some(pass) = url.password() {
Some(
percent_decode(pass.as_ref())
.decode_utf8_lossy()
.into_owned(),
)
} else {
None
}
}
fn get_opts_db_name_from_url(url: &Url) -> Option<String> {
if let Some(mut segments) = url.path_segments() {
segments
.next()
.filter(|&db_name| !db_name.is_empty())
.map(|db_name| {
percent_decode(db_name.as_ref())
.decode_utf8_lossy()
.into_owned()
})
} else {
None
}
}
fn from_url_basic(url_str: &str) -> Result<(Opts, Vec<(String, String)>), UrlError> {
let url = Url::parse(url_str)?;
if url.scheme() != "mysql" {
return Err(UrlError::UnsupportedScheme(url.scheme().to_string()));
}
if url.cannot_be_a_base() {
return Err(UrlError::BadUrl);
}
let user = get_opts_user_from_url(&url);
let pass = get_opts_pass_from_url(&url);
let ip_or_hostname = url
.host()
.ok_or(UrlError::BadUrl)
.and_then(|host| url::Host::parse(&host.to_string()).map_err(|_| UrlError::BadUrl))?;
let tcp_port = url.port().unwrap_or(3306);
let db_name = get_opts_db_name_from_url(&url);
let query_pairs = url.query_pairs().into_owned().collect();
let opts = Opts(Box::new(InnerOpts {
user,
pass,
ip_or_hostname,
tcp_port,
db_name,
..InnerOpts::default()
}));
Ok((opts, query_pairs))
}
fn from_url(url: &str) -> Result<Opts, UrlError> {
let (mut opts, query_pairs) = from_url_basic(url)?;
for (key, value) in query_pairs {
if key == "prefer_socket" {
if value == "true" {
opts.set_prefer_socket(true);
} else if value == "false" {
opts.set_prefer_socket(false);
} else {
return Err(UrlError::InvalidValue("prefer_socket".into(), value));
}
} else if key == "tcp_keepalive_time_ms" {
match u32::from_str(&*value) {
Ok(tcp_keepalive_time_ms) => {
opts.0.tcp_keepalive_time = Some(tcp_keepalive_time_ms);
}
_ => {
return Err(UrlError::InvalidValue(
"tcp_keepalive_time_ms".into(),
value,
));
}
}
} else if key == "tcp_connect_timeout_ms" {
match u64::from_str(&*value) {
Ok(tcp_connect_timeout_ms) => {
opts.0.tcp_connect_timeout =
Some(Duration::from_millis(tcp_connect_timeout_ms));
}
_ => {
return Err(UrlError::InvalidValue(
"tcp_connect_timeout_ms".into(),
value,
));
}
}
} else if key == "stmt_cache_size" {
match usize::from_str(&*value) {
Ok(stmt_cache_size) => {
opts.0.stmt_cache_size = stmt_cache_size;
}
_ => {
return Err(UrlError::InvalidValue("stmt_cache_size".into(), value));
}
}
} else if key == "compress" {
if value == "true" {
opts.0.compress = Some(crate::Compression::default());
} else if value == "fast" {
opts.0.compress = Some(crate::Compression::fast());
} else if value == "best" {
opts.0.compress = Some(crate::Compression::best());
} else if value.len() == 1 && 0x30 <= value.as_bytes()[0] && value.as_bytes()[0] <= 0x39
{
opts.0.compress =
Some(crate::Compression::new((value.as_bytes()[0] - 0x30) as u32));
} else {
return Err(UrlError::InvalidValue("compress".into(), value));
}
} else if key == "socket" {
let socket = percent_decode_str(&*value).decode_utf8_lossy().into_owned();
if !socket.is_empty() {
opts.0.socket = Some(socket);
}
} else {
return Err(UrlError::UnknownParameter(key));
}
}
Ok(opts)
}
impl<S: AsRef<str>> From<S> for Opts {
fn from(url: S) -> Opts {
match from_url(url.as_ref()) {
Ok(opts) => opts,
Err(err) => panic!("{}", err),
}
}
}
#[cfg(test)]
mod test {
use mysql_common::proto::codec::Compression;
use super::{InnerOpts, Opts};
#[test]
fn should_report_empty_url_database_as_none() {
let opt = Opts::from("mysql://localhost/");
assert_eq!(opt.get_db_name(), None);
}
#[test]
fn should_convert_url_into_opts() {
let opts = "mysql://us%20r:p%20w@localhost:3308/db%2dname?prefer_socket=false&tcp_keepalive_time_ms=5000&socket=%2Ftmp%2Fmysql.sock&compress=8";
assert_eq!(
Opts(Box::new(InnerOpts {
user: Some("us r".to_string()),
pass: Some("p w".to_string()),
ip_or_hostname: url::Host::Domain("localhost".to_string()),
tcp_port: 3308,
db_name: Some("db-name".to_string()),
prefer_socket: false,
tcp_keepalive_time: Some(5000),
socket: Some("/tmp/mysql.sock".into()),
compress: Some(Compression::new(8)),
..InnerOpts::default()
})),
opts.into()
);
}
#[test]
#[should_panic]
fn should_panic_on_invalid_url() {
let opts = "42";
let _: Opts = opts.into();
}
#[test]
#[should_panic]
fn should_panic_on_invalid_scheme() {
let opts = "postgres://localhost";
let _: Opts = opts.into();
}
#[test]
#[should_panic]
fn should_panic_on_unknown_query_param() {
let opts = "mysql://localhost/foo?bar=baz";
let _: Opts = opts.into();
}
}