From bf1ac7f3681fd058f755ff45f67fb280c2450b2f Mon Sep 17 00:00:00 2001 From: Bassem Girgis Date: Fri, 20 Aug 2021 01:17:52 -0500 Subject: [PATCH] improve tests --- Cargo.toml | 5 + db/mssql/init.sql | 20 +++- src/check_maria_db.rs | 200 +++++++++++++++++++++++++++++++++++++++- src/check_ms_sql.rs | 197 ++++++++++++++++++++++++++++++++++++++-- src/check_mysql.rs | 200 +++++++++++++++++++++++++++++++++++++++- src/check_postgres.rs | 207 +++++++++++++++++++++++++++++++++++++++++- 6 files changed, 809 insertions(+), 20 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 9c53d1d..da11d8a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -12,6 +12,11 @@ log = { version = "^0.4", features = [ ] } fern = "^0.5" +serde = { version = "^1.0.127", features = ["derive"] } +serde_json = "^1.0.66" + +uuid = { version = "0.8.2", features = ["v4"] } + # sqlx sqlx = { version = "0.5.6", features = [ # tokio + rustls diff --git a/db/mssql/init.sql b/db/mssql/init.sql index 9988a30..29153a2 100644 --- a/db/mssql/init.sql +++ b/db/mssql/init.sql @@ -18,10 +18,20 @@ GO CREATE USER demo_user WITH PASSWORD = 'pas$w0rd'; GO -CREATE TABLE Products (ID int, ProductName nvarchar(max)); +-- a single table is used for all events in the cqrs system +CREATE TABLE events +( + aggregate_type VARCHAR(256) NOT NULL, + aggregate_id VARCHAR(256) NOT NULL, + sequence bigint CHECK (sequence >= 0) NOT NULL, + payload TEXT NOT NULL, + metadata TEXT NOT NULL, + PRIMARY KEY (aggregate_type, aggregate_id, sequence) +); + GO -GRANT SELECT ON OBJECT::dbo.Products TO demo_user; -GRANT INSERT ON OBJECT::dbo.Products TO demo_user; -GRANT UPDATE ON OBJECT::dbo.Products TO demo_user; -GRANT DELETE ON OBJECT::dbo.Products TO demo_user; +GRANT SELECT ON OBJECT::dbo.events TO demo_user; +GRANT INSERT ON OBJECT::dbo.events TO demo_user; +GRANT UPDATE ON OBJECT::dbo.events TO demo_user; +GRANT DELETE ON OBJECT::dbo.events TO demo_user; diff --git a/src/check_maria_db.rs b/src/check_maria_db.rs index 3af8077..9c41ee9 100644 --- a/src/check_maria_db.rs +++ b/src/check_maria_db.rs @@ -1,11 +1,63 @@ use log::info; +use std::collections::HashMap; + use sqlx::mysql::{ MySqlConnectOptions, + MySqlPool, MySqlPoolOptions, }; -pub async fn check_maria_db() -> Result<(), sqlx::Error> { +pub static INSERT_EVENT: &str = " +INSERT INTO + events + ( + aggregate_type, + aggregate_id, + sequence, + payload, + metadata + ) +VALUES + ( + ?, + ?, + ?, + ?, + ? + ); +"; + +pub static SELECT_EVENTS_WITH_METADATA: &str = " +SELECT + sequence, + payload, + metadata +FROM + events +WHERE + aggregate_type = ? + AND + aggregate_id = ? +ORDER BY + sequence; +"; + +pub static UPDATE_EVENTS: &str = " +UPDATE + events +SET + payload = ?, + metadata = ? +WHERE + aggregate_type = ? + AND + aggregate_id = ? + AND + sequence = ?; +"; + +pub async fn get_pool() -> Result { // "mysql://demo_user:demo_pass@localhost:8084/demo" let options = MySqlConnectOptions::new() .host("localhost") @@ -19,12 +71,152 @@ pub async fn check_maria_db() -> Result<(), sqlx::Error> { .connect_with(options) .await?; - // Make a simple query to return the given parameter - let row = sqlx::query("SELECT * from events") + 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!("Received {:?}", row); + info!("select success"); + + for row in rows { + let payload: HashMap = + serde_json::from_value(row.1).unwrap(); + let metadata: HashMap = + 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 = + serde_json::from_value(row.1).unwrap(); + let metadata: HashMap = + 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(()) } diff --git a/src/check_ms_sql.rs b/src/check_ms_sql.rs index 9542755..f94f32d 100644 --- a/src/check_ms_sql.rs +++ b/src/check_ms_sql.rs @@ -1,11 +1,62 @@ use log::info; +use std::collections::HashMap; + use sqlx::mssql::{ MssqlConnectOptions, MssqlPool, }; -pub async fn check_ms_sql() -> Result<(), sqlx::Error> { +pub static INSERT_EVENT: &str = " +INSERT INTO + events + ( + aggregate_type, + aggregate_id, + sequence, + payload, + metadata + ) +VALUES + ( + ?, + ?, + ?, + ?, + ? + ); +"; + +pub static SELECT_EVENTS_WITH_METADATA: &str = " +SELECT + sequence, + payload, + metadata +FROM + events +WHERE + aggregate_type = $1 + AND + aggregate_id = $2 +ORDER BY + sequence; +"; + +pub static UPDATE_EVENTS: &str = " +UPDATE + events +SET + payload = $4, + metadata = $5 +WHERE + aggregate_type = $1 + AND + aggregate_id = $2 + AND + sequence = $3; +"; + +pub async fn get_pool() -> Result { // "mssql://sa:'adm1n_pa%s'@localhost:8085/demo" let options = MssqlConnectOptions::new() .host("localhost") @@ -16,12 +67,146 @@ pub async fn check_ms_sql() -> Result<(), sqlx::Error> { let pool = MssqlPool::connect_with(options).await?; - // Make a simple query to return the given parameter - let _row = sqlx::query("SELECT * from Products") - .fetch_all(&pool) - .await?; + Ok(pool) +} - info!("Received ok!"); +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_string(&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_string(&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, String, String)> = + 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 = + serde_json::from_str(&row.1).unwrap(); + let metadata: HashMap = + serde_json::from_str(&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_string(&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_string(&metadata) { + Ok(x) => x, + Err(e) => { + panic!( + "metadata serialization error '{}'", + e.to_string() + ); + }, + }; + + let rows_affected = sqlx::query(UPDATE_EVENTS) + .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, String, String)> = + 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 = + serde_json::from_str(&row.1).unwrap(); + let metadata: HashMap = + serde_json::from_str(&row.2).unwrap(); + info!( + "{},{:?},{:?}", + &row.0, payload, metadata + ) + } + + Ok(()) +} + +pub async fn check_ms_sql() -> Result<(), sqlx::Error> { + test_insert_select_update().await?; Ok(()) } diff --git a/src/check_mysql.rs b/src/check_mysql.rs index a06e826..e709adf 100644 --- a/src/check_mysql.rs +++ b/src/check_mysql.rs @@ -1,11 +1,63 @@ use log::info; +use std::collections::HashMap; + use sqlx::mysql::{ MySqlConnectOptions, + MySqlPool, MySqlPoolOptions, }; -pub async fn check_mysql() -> Result<(), sqlx::Error> { +pub static INSERT_EVENT: &str = " +INSERT INTO + events + ( + aggregate_type, + aggregate_id, + sequence, + payload, + metadata + ) +VALUES + ( + ?, + ?, + ?, + ?, + ? + ); +"; + +pub static SELECT_EVENTS_WITH_METADATA: &str = " +SELECT + sequence, + payload, + metadata +FROM + events +WHERE + aggregate_type = ? + AND + aggregate_id = ? +ORDER BY + sequence; +"; + +pub static UPDATE_EVENTS: &str = " +UPDATE + events +SET + payload = ?, + metadata = ? +WHERE + aggregate_type = ? + AND + aggregate_id = ? + AND + sequence = ?; +"; + +pub async fn get_pool() -> Result { // "mysql://demo_user:demo_pass@localhost:8086/demo" let options = MySqlConnectOptions::new() .host("localhost") @@ -19,12 +71,152 @@ pub async fn check_mysql() -> Result<(), sqlx::Error> { .connect_with(options) .await?; - // Make a simple query to return the given parameter - let row = sqlx::query("SELECT * from events") + 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!("Received {:?}", row); + info!("select success"); + + for row in rows { + let payload: HashMap = + serde_json::from_value(row.1).unwrap(); + let metadata: HashMap = + 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 = + serde_json::from_value(row.1).unwrap(); + let metadata: HashMap = + serde_json::from_value(row.2).unwrap(); + info!( + "{},{:?},{:?}", + &row.0, payload, metadata + ) + } + + Ok(()) +} + +pub async fn check_mysql() -> Result<(), sqlx::Error> { + test_insert_select_update().await?; Ok(()) } diff --git a/src/check_postgres.rs b/src/check_postgres.rs index 0ee6a01..db386e9 100644 --- a/src/check_postgres.rs +++ b/src/check_postgres.rs @@ -1,11 +1,63 @@ use log::info; +use std::collections::HashMap; + use sqlx::postgres::{ PgConnectOptions, + PgPool, PgPoolOptions, }; -pub async fn check_postgres() -> Result<(), sqlx::Error> { +pub static INSERT_EVENT: &str = " +INSERT INTO + events + ( + aggregate_type, + aggregate_id, + sequence, + payload, + metadata + ) +VALUES + ( + $1, + $2, + $3, + $4, + $5 + ); +"; + +pub static SELECT_EVENTS_WITH_METADATA: &str = " +SELECT + sequence, + payload, + metadata +FROM + events +WHERE + aggregate_type = $1 + AND + aggregate_id = $2 +ORDER BY + sequence; +"; + +pub static UPDATE_EVENTS: &str = " +UPDATE + events +SET + payload = $4, + metadata = $5 +WHERE + aggregate_type = $1 + AND + aggregate_id = $2 + AND + sequence = $3; +"; + +async fn get_pool() -> Result { // "postgres://demo_user:demo_pass@localhost:8087/demo" let options = PgConnectOptions::new() .host("localhost") @@ -19,6 +71,12 @@ pub async fn check_postgres() -> Result<(), sqlx::Error> { .connect_with(options) .await?; + Ok(pool) +} + +async fn test_connect() -> Result<(), sqlx::Error> { + let pool = get_pool().await?; + // Make a simple query to return the given parameter let row: (i64,) = sqlx::query_as("SELECT $1") .bind(150_i64) @@ -28,6 +86,153 @@ pub async fn check_postgres() -> Result<(), sqlx::Error> { assert_eq!(row.0, 150); info!("Received {}", row.0); + Ok(()) +} + +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 = + serde_json::from_value(row.1).unwrap(); + let metadata: HashMap = + 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(&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 = + serde_json::from_value(row.1).unwrap(); + let metadata: HashMap = + serde_json::from_value(row.2).unwrap(); + info!( + "{},{:?},{:?}", + &row.0, payload, metadata + ) + } + + Ok(()) +} + +pub async fn check_postgres() -> Result<(), sqlx::Error> { + test_connect().await?; + test_insert_select_update().await?; Ok(()) }