Files
test-sqlx/src/check_maria_db.rs
2021-08-20 15:05:50 -05:00

223 lines
4.6 KiB
Rust

use log::info;
use std::collections::HashMap;
use sqlx::mysql::{
MySqlConnectOptions,
MySqlPool,
MySqlPoolOptions,
};
static INSERT_EVENT: &str = "
INSERT INTO
events
(
aggregate_type,
aggregate_id,
sequence,
payload,
metadata
)
VALUES
(
?,
?,
?,
?,
?
);
";
static SELECT_EVENTS_WITH_METADATA: &str = "
SELECT
sequence,
payload,
metadata
FROM
events
WHERE
aggregate_type = ?
AND
aggregate_id = ?
ORDER BY
sequence;
";
static UPDATE_EVENTS: &str = "
UPDATE
events
SET
payload = ?,
metadata = ?
WHERE
aggregate_type = ?
AND
aggregate_id = ?
AND
sequence = ?;
";
pub async fn get_pool() -> Result<MySqlPool, sqlx::Error> {
// "mysql://demo_user:demo_pass@localhost:8081/demo"
let options = MySqlConnectOptions::new()
.host("localhost")
.port(8081)
.database("demo")
.username("demo_user")
.password("demo_pass");
let pool = MySqlPoolOptions::new()
.max_connections(5)
.connect_with(options)
.await?;
Ok(pool)
}
pub async fn test_insert_select_update() -> Result<(), sqlx::Error> {
let pool = get_pool().await?;
// insert
let mut payload = HashMap::new();
payload.insert("k".to_string(), "v".to_string());
let payload = match serde_json::to_value(payload) {
Ok(x) => x,
Err(e) => {
panic!(
"payload serialization error '{}'",
e.to_string()
);
},
};
let mut metadata = HashMap::new();
metadata.insert("k1".to_string(), "v1".to_string());
let metadata = match serde_json::to_value(metadata) {
Ok(x) => x,
Err(e) => {
panic!(
"metadata serialization error '{}'",
e.to_string()
);
},
};
let aggregate_type = uuid::Uuid::new_v4().to_string();
let aggregate_id = uuid::Uuid::new_v4().to_string();
let rows_affected = sqlx::query(INSERT_EVENT)
.bind(&aggregate_type)
.bind(&aggregate_id)
.bind(1)
.bind(payload)
.bind(metadata)
.execute(&pool)
.await?
.rows_affected();
info!(
"Insert affected '{}' rows",
rows_affected
);
// select
let rows: Vec<(
i64,
serde_json::Value,
serde_json::Value,
)> = sqlx::query_as(SELECT_EVENTS_WITH_METADATA)
.bind(&aggregate_type)
.bind(&aggregate_id)
.fetch_all(&pool)
.await?;
info!("select success");
for row in rows {
let payload: HashMap<String, String> =
serde_json::from_value(row.1).unwrap();
let metadata: HashMap<String, String> =
serde_json::from_value(row.2).unwrap();
info!(
"{},{:?},{:?}",
&row.0, payload, metadata
)
}
// update
let mut payload = HashMap::new();
payload.insert("k2".to_string(), "v2".to_string());
let payload = match serde_json::to_value(payload) {
Ok(x) => x,
Err(e) => {
panic!(
"payload serialization error '{}'",
e.to_string()
);
},
};
let mut metadata = HashMap::new();
metadata.insert("k3".to_string(), "v3".to_string());
let metadata = match serde_json::to_value(metadata) {
Ok(x) => x,
Err(e) => {
panic!(
"metadata serialization error '{}'",
e.to_string()
);
},
};
let rows_affected = sqlx::query(UPDATE_EVENTS)
.bind(payload)
.bind(metadata)
.bind(&aggregate_type)
.bind(&aggregate_id)
.bind(1)
.execute(&pool)
.await?
.rows_affected();
info!(
"Insert affected '{}' rows",
rows_affected
);
// select
let rows: Vec<(
i64,
serde_json::Value,
serde_json::Value,
)> = sqlx::query_as(SELECT_EVENTS_WITH_METADATA)
.bind(&aggregate_type)
.bind(&aggregate_id)
.fetch_all(&pool)
.await?;
info!("select success");
for row in rows {
let payload: HashMap<String, String> =
serde_json::from_value(row.1).unwrap();
let metadata: HashMap<String, String> =
serde_json::from_value(row.2).unwrap();
info!(
"{},{:?},{:?}",
&row.0, payload, metadata
)
}
Ok(())
}
pub async fn check_maria_db() -> Result<(), sqlx::Error> {
test_insert_select_update().await?;
Ok(())
}