diff --git a/Makefile b/Makefile index 48f6eea..68ab806 100644 --- a/Makefile +++ b/Makefile @@ -5,6 +5,7 @@ all: clean: rm -rf target + rm -rf test.db* build: cargo build @@ -14,6 +15,7 @@ test: up: docker-compose up -d + rm -rf test.db* down: docker-compose down diff --git a/src/check_maria_db.rs b/src/check_maria_db.rs index 9c41ee9..464881f 100644 --- a/src/check_maria_db.rs +++ b/src/check_maria_db.rs @@ -8,7 +8,7 @@ use sqlx::mysql::{ MySqlPoolOptions, }; -pub static INSERT_EVENT: &str = " +static INSERT_EVENT: &str = " INSERT INTO events ( @@ -28,7 +28,7 @@ VALUES ); "; -pub static SELECT_EVENTS_WITH_METADATA: &str = " +static SELECT_EVENTS_WITH_METADATA: &str = " SELECT sequence, payload, @@ -43,7 +43,7 @@ ORDER BY sequence; "; -pub static UPDATE_EVENTS: &str = " +static UPDATE_EVENTS: &str = " UPDATE events SET diff --git a/src/check_ms_sql.rs b/src/check_ms_sql.rs index f94f32d..40cff7b 100644 --- a/src/check_ms_sql.rs +++ b/src/check_ms_sql.rs @@ -7,7 +7,7 @@ use sqlx::mssql::{ MssqlPool, }; -pub static INSERT_EVENT: &str = " +static INSERT_EVENT: &str = " INSERT INTO events ( @@ -27,7 +27,7 @@ VALUES ); "; -pub static SELECT_EVENTS_WITH_METADATA: &str = " +static SELECT_EVENTS_WITH_METADATA: &str = " SELECT sequence, payload, @@ -42,7 +42,7 @@ ORDER BY sequence; "; -pub static UPDATE_EVENTS: &str = " +static UPDATE_EVENTS: &str = " UPDATE events SET diff --git a/src/check_mysql.rs b/src/check_mysql.rs index e709adf..0b6f64a 100644 --- a/src/check_mysql.rs +++ b/src/check_mysql.rs @@ -8,7 +8,7 @@ use sqlx::mysql::{ MySqlPoolOptions, }; -pub static INSERT_EVENT: &str = " +static INSERT_EVENT: &str = " INSERT INTO events ( @@ -28,7 +28,7 @@ VALUES ); "; -pub static SELECT_EVENTS_WITH_METADATA: &str = " +static SELECT_EVENTS_WITH_METADATA: &str = " SELECT sequence, payload, @@ -43,7 +43,7 @@ ORDER BY sequence; "; -pub static UPDATE_EVENTS: &str = " +static UPDATE_EVENTS: &str = " UPDATE events SET diff --git a/src/check_postgres.rs b/src/check_postgres.rs index db386e9..722c7f3 100644 --- a/src/check_postgres.rs +++ b/src/check_postgres.rs @@ -8,7 +8,7 @@ use sqlx::postgres::{ PgPoolOptions, }; -pub static INSERT_EVENT: &str = " +static INSERT_EVENT: &str = " INSERT INTO events ( @@ -28,7 +28,7 @@ VALUES ); "; -pub static SELECT_EVENTS_WITH_METADATA: &str = " +static SELECT_EVENTS_WITH_METADATA: &str = " SELECT sequence, payload, @@ -43,7 +43,7 @@ ORDER BY sequence; "; -pub static UPDATE_EVENTS: &str = " +static UPDATE_EVENTS: &str = " UPDATE events SET diff --git a/src/check_sqlite.rs b/src/check_sqlite.rs index d2ba72a..1496fa5 100644 --- a/src/check_sqlite.rs +++ b/src/check_sqlite.rs @@ -1,11 +1,77 @@ use log::info; +use std::collections::HashMap; + use sqlx::sqlite::{ SqliteConnectOptions, + SqlitePool, SqlitePoolOptions, }; -pub async fn check_sqlite() -> Result<(), sqlx::Error> { +static CREATE_TABLE: &str = " +CREATE TABLE IF NOT EXISTS +events + ( + aggregate_type TEXT NOT NULL, + aggregate_id TEXT NOT NULL, + sequence bigint CHECK (sequence >= 0) NOT NULL, + payload TEXT NOT NULL, + metadata TEXT NOT NULL, + timestamp timestamp DEFAULT (CURRENT_TIMESTAMP), + PRIMARY KEY (aggregate_type, aggregate_id, sequence) + ); +"; + +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 { // "sqlite://test.db" let options = SqliteConnectOptions::new() .filename("test.db") @@ -16,6 +82,12 @@ pub async fn check_sqlite() -> Result<(), sqlx::Error> { .connect_with(options) .await?; + Ok(pool) +} + +pub 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,3 +100,162 @@ pub async fn check_sqlite() -> Result<(), sqlx::Error> { Ok(()) } + +async fn test_insert_select_update() -> Result<(), sqlx::Error> { + let pool = get_pool().await?; + + // create + let rows_affected = sqlx::query(CREATE_TABLE) + .execute(&pool) + .await? + .rows_affected(); + + info!( + "Create affected '{}' rows", + rows_affected + ); + + // 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(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_sqlite() -> Result<(), sqlx::Error> { + test_connect().await?; + test_insert_select_update().await?; + + Ok(()) +} diff --git a/src/main.rs b/src/main.rs index bb41d9e..e691d6e 100644 --- a/src/main.rs +++ b/src/main.rs @@ -37,19 +37,44 @@ async fn main() -> Result<(), sqlx::Error> { let ts = vec![ tokio::spawn(async move { - check_postgres().await.unwrap(); + match check_maria_db().await { + Ok(()) => {}, + Err(e) => { + panic!("MARIADB ERROR '{}'", e.to_string()) + }, + }; + }), + // tokio::spawn(async move { + // match check_ms_sql().await { + // Ok(()) => {}, + // Err(e) => { + // panic!("MSSQL ERROR '{}'", e.to_string()) + // }, + // }; + // }), + tokio::spawn(async move { + match check_mysql().await { + Ok(()) => {}, + Err(e) => { + panic!("MYSQL ERROR '{}'", e.to_string()) + }, + }; }), tokio::spawn(async move { - check_mysql().await.unwrap(); + match check_postgres().await { + Ok(()) => {}, + Err(e) => { + panic!("POSTGRES ERROR '{}'", e.to_string()) + }, + }; }), tokio::spawn(async move { - check_ms_sql().await.unwrap(); - }), - tokio::spawn(async move { - check_maria_db().await.unwrap(); - }), - tokio::spawn(async move { - check_sqlite().await.unwrap(); + match check_sqlite().await { + Ok(()) => {}, + Err(e) => { + panic!("SQLITE ERROR '{}'", e.to_string()) + }, + }; }), ];