Compare commits

...

10 Commits

Author SHA1 Message Date
a297eb285f minor edit 2021-08-25 17:15:49 -05:00
ff2d7b6504 minor 2021-08-20 15:23:18 -05:00
f9f090605c minor 2021-08-20 15:13:48 -05:00
d9691c0cc9 minor 2021-08-20 15:08:23 -05:00
01d2afb3dd cleanup minor things 2021-08-20 15:05:50 -05:00
4b859492c1 minor 2021-08-20 01:58:02 -05:00
80916d3a49 minor 2021-08-20 01:50:19 -05:00
9723f350f9 minor changes 2021-08-20 01:43:45 -05:00
e49f0c90ce add sqlite 2021-08-20 01:33:20 -05:00
bf1ac7f368 improve tests 2021-08-20 01:17:52 -05:00
14 changed files with 1129 additions and 57 deletions

4
.gitignore vendored
View File

@@ -1,4 +1,6 @@
/target /target
Cargo.lock Cargo.lock
*.log *.log
test.db*
demo.db*

View File

@@ -12,6 +12,13 @@ log = { version = "^0.4", features = [
] } ] }
fern = "^0.5" fern = "^0.5"
# serialization
serde = { version = "^1.0.127", features = ["derive"] }
serde_json = "^1.0.66"
# ids
uuid = { version = "0.8.2", features = ["v4"] }
# sqlx # sqlx
sqlx = { version = "0.5.6", features = [ sqlx = { version = "0.5.6", features = [
# tokio + rustls # tokio + rustls

View File

@@ -5,6 +5,9 @@ all:
clean: clean:
rm -rf target rm -rf target
rm -rf Cargo.lock
rm -rf *.log
rm -rf demo.db*
build: build:
cargo build cargo build
@@ -14,6 +17,7 @@ test:
up: up:
docker-compose up -d docker-compose up -d
rm -rf demo.db*
down: down:
docker-compose down docker-compose down

View File

@@ -1 +1,21 @@
# sqlx Test # `sqlx` Test
This is a test program to try few things with the `sqlx` library.
The idea is to figure out everything needed to support
moving [cqrs-es2](https://github.com/brgirgis/cqrs-es2) to use `sqlx`.
## Build
To build the executable simply trigger `Cargo`
cargo build
## Usage
To run the executable you will need to spin up the database stack:
docker-compose up -d
Wait for a few seconds for the stack to be ready and then you can run the executable:
cargo run

View File

@@ -4,11 +4,11 @@ USE demo;
-- a single table is used for all events in the cqrs system -- a single table is used for all events in the cqrs system
CREATE TABLE events CREATE TABLE events
( (
aggregate_type VARCHAR(256) NOT NULL, aggregate_type VARCHAR(256) NOT NULL,
aggregate_id VARCHAR(256) NOT NULL, aggregate_id VARCHAR(256) NOT NULL,
sequence bigint CHECK (sequence >= 0) , sequence bigint CHECK (sequence >= 0) ,
payload TEXT , payload TEXT ,
metadata TEXT , metadata TEXT ,
timestamp timestamp DEFAULT (CURRENT_TIMESTAMP), timestamp timestamp DEFAULT (CURRENT_TIMESTAMP),
PRIMARY KEY (aggregate_type, aggregate_id, sequence) PRIMARY KEY (aggregate_type, aggregate_id, sequence)
); );

View File

@@ -18,10 +18,20 @@ GO
CREATE USER demo_user WITH PASSWORD = 'pas$w0rd'; CREATE USER demo_user WITH PASSWORD = 'pas$w0rd';
GO 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 GO
GRANT SELECT ON OBJECT::dbo.Products TO demo_user; GRANT SELECT ON OBJECT::dbo.events TO demo_user;
GRANT INSERT ON OBJECT::dbo.Products TO demo_user; GRANT INSERT ON OBJECT::dbo.events TO demo_user;
GRANT UPDATE ON OBJECT::dbo.Products TO demo_user; GRANT UPDATE ON OBJECT::dbo.events TO demo_user;
GRANT DELETE ON OBJECT::dbo.Products TO demo_user; GRANT DELETE ON OBJECT::dbo.events TO demo_user;

View File

@@ -4,11 +4,11 @@ USE demo;
-- a single table is used for all events in the cqrs system -- a single table is used for all events in the cqrs system
CREATE TABLE events CREATE TABLE events
( (
aggregate_type VARCHAR(256) NOT NULL, aggregate_type VARCHAR(256) NOT NULL,
aggregate_id VARCHAR(256) NOT NULL, aggregate_id VARCHAR(256) NOT NULL,
sequence bigint CHECK (sequence >= 0) NOT NULL, sequence bigint CHECK (sequence >= 0) NOT NULL,
payload TEXT NOT NULL, payload TEXT NOT NULL,
metadata TEXT NOT NULL, metadata TEXT NOT NULL,
timestamp timestamp DEFAULT (CURRENT_TIMESTAMP), timestamp timestamp DEFAULT (CURRENT_TIMESTAMP),
PRIMARY KEY (aggregate_type, aggregate_id, sequence) PRIMARY KEY (aggregate_type, aggregate_id, sequence)
); );

View File

@@ -8,7 +8,7 @@ services:
networks: networks:
- default - default
ports: ports:
- "8084:3306" - "8081:3306"
environment: environment:
#- "MARIADB_USER=root" #- "MARIADB_USER=root"
- "MARIADB_ROOT_PASSWORD=admin_pass" - "MARIADB_ROOT_PASSWORD=admin_pass"
@@ -24,7 +24,7 @@ services:
networks: networks:
- default - default
ports: ports:
- "8085:1433" - "8082:1433"
environment: environment:
#- "SA_USER=sa" #- "SA_USER=sa"
- "SA_PASSWORD=adm1n_pa%s" - "SA_PASSWORD=adm1n_pa%s"
@@ -37,7 +37,7 @@ services:
networks: networks:
- default - default
ports: ports:
- "8086:3306" - "8083:3306"
environment: environment:
#- "MYSQL_USER=root" #- "MYSQL_USER=root"
- "MYSQL_ROOT_PASSWORD=admin_pass" - "MYSQL_ROOT_PASSWORD=admin_pass"
@@ -50,7 +50,7 @@ services:
networks: networks:
- default - default
ports: ports:
- "8087:5432" - "8084:5432"
environment: environment:
- "POSTGRES_USER=admin" - "POSTGRES_USER=admin"
- "POSTGRES_PASSWORD=admin_pass" - "POSTGRES_PASSWORD=admin_pass"
@@ -61,9 +61,8 @@ services:
networks: networks:
- default - default
ports: ports:
- "8083:8080" - "8080:8080"
depends_on: depends_on:
- maria-db - maria-db
- mssql-db
- mysql-db - mysql-db
- postgres-db - postgres-db

View File

@@ -1,15 +1,67 @@
use log::info; use log::info;
use std::collections::HashMap;
use sqlx::mysql::{ use sqlx::mysql::{
MySqlConnectOptions, MySqlConnectOptions,
MySqlPool,
MySqlPoolOptions, MySqlPoolOptions,
}; };
pub async fn check_maria_db() -> Result<(), sqlx::Error> { static INSERT_EVENT: &str = "
// "mysql://demo_user:demo_pass@localhost:8084/demo" 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() let options = MySqlConnectOptions::new()
.host("localhost") .host("localhost")
.port(8084) .port(8081)
.database("demo") .database("demo")
.username("demo_user") .username("demo_user")
.password("demo_pass"); .password("demo_pass");
@@ -19,12 +71,152 @@ pub async fn check_maria_db() -> Result<(), sqlx::Error> {
.connect_with(options) .connect_with(options)
.await?; .await?;
// Make a simple query to return the given parameter Ok(pool)
let row = sqlx::query("SELECT * from events") }
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) .fetch_all(&pool)
.await?; .await?;
info!("Received {:?}", row); 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(()) Ok(())
} }

View File

@@ -1,27 +1,212 @@
use log::info; use log::info;
use std::collections::HashMap;
use sqlx::mssql::{ use sqlx::mssql::{
MssqlConnectOptions, MssqlConnectOptions,
MssqlPool, MssqlPool,
}; };
pub async fn check_ms_sql() -> Result<(), sqlx::Error> { static INSERT_EVENT: &str = "
// "mssql://sa:'adm1n_pa%s'@localhost:8085/demo" 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 = $1
AND
aggregate_id = $2
ORDER BY
sequence;
";
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<MssqlPool, sqlx::Error> {
// "mssql://sa:'adm1n_pa%s'@localhost:8082/demo"
let options = MssqlConnectOptions::new() let options = MssqlConnectOptions::new()
.host("localhost") .host("localhost")
.port(8085) .port(8082)
.database("demo") .database("demo")
.username("demo_user") .username("demo_user")
.password("pas$w0rd"); .password("pas$w0rd");
let pool = MssqlPool::connect_with(options).await?; let pool = MssqlPool::connect_with(options).await?;
// Make a simple query to return the given parameter Ok(pool)
let _row = sqlx::query("SELECT * from Products") }
.fetch_all(&pool)
.await?;
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<String, String> =
serde_json::from_str(&row.1).unwrap();
let metadata: HashMap<String, String> =
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<String, String> =
serde_json::from_str(&row.1).unwrap();
let metadata: HashMap<String, String> =
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(()) Ok(())
} }

View File

@@ -1,15 +1,67 @@
use log::info; use log::info;
use std::collections::HashMap;
use sqlx::mysql::{ use sqlx::mysql::{
MySqlConnectOptions, MySqlConnectOptions,
MySqlPool,
MySqlPoolOptions, MySqlPoolOptions,
}; };
pub async fn check_mysql() -> Result<(), sqlx::Error> { static INSERT_EVENT: &str = "
// "mysql://demo_user:demo_pass@localhost:8086/demo" 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:8083/demo"
let options = MySqlConnectOptions::new() let options = MySqlConnectOptions::new()
.host("localhost") .host("localhost")
.port(8086) .port(8083)
.database("demo") .database("demo")
.username("demo_user") .username("demo_user")
.password("demo_pass"); .password("demo_pass");
@@ -19,12 +71,152 @@ pub async fn check_mysql() -> Result<(), sqlx::Error> {
.connect_with(options) .connect_with(options)
.await?; .await?;
// Make a simple query to return the given parameter Ok(pool)
let row = sqlx::query("SELECT * from events") }
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) .fetch_all(&pool)
.await?; .await?;
info!("Received {:?}", row); 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_mysql() -> Result<(), sqlx::Error> {
test_insert_select_update().await?;
Ok(()) Ok(())
} }

View File

@@ -1,15 +1,67 @@
use log::info; use log::info;
use std::collections::HashMap;
use sqlx::postgres::{ use sqlx::postgres::{
PgConnectOptions, PgConnectOptions,
PgPool,
PgPoolOptions, PgPoolOptions,
}; };
pub async fn check_postgres() -> Result<(), sqlx::Error> { static INSERT_EVENT: &str = "
// "postgres://demo_user:demo_pass@localhost:8087/demo" INSERT INTO
events
(
aggregate_type,
aggregate_id,
sequence,
payload,
metadata
)
VALUES
(
$1,
$2,
$3,
$4,
$5
);
";
static SELECT_EVENTS_WITH_METADATA: &str = "
SELECT
sequence,
payload,
metadata
FROM
events
WHERE
aggregate_type = $1
AND
aggregate_id = $2
ORDER BY
sequence;
";
static UPDATE_EVENTS: &str = "
UPDATE
events
SET
payload = $1,
metadata = $2
WHERE
aggregate_type = $3
AND
aggregate_id = $4
AND
sequence = $5;
";
async fn get_pool() -> Result<PgPool, sqlx::Error> {
// "postgres://demo_user:demo_pass@localhost:8084/demo"
let options = PgConnectOptions::new() let options = PgConnectOptions::new()
.host("localhost") .host("localhost")
.port(8087) .port(8084)
.database("demo") .database("demo")
.username("demo_user") .username("demo_user")
.password("demo_pass"); .password("demo_pass");
@@ -19,6 +71,12 @@ pub async fn check_postgres() -> Result<(), sqlx::Error> {
.connect_with(options) .connect_with(options)
.await?; .await?;
Ok(pool)
}
async fn test_connect() -> Result<(), sqlx::Error> {
let pool = get_pool().await?;
// Make a simple query to return the given parameter // Make a simple query to return the given parameter
let row: (i64,) = sqlx::query_as("SELECT $1") let row: (i64,) = sqlx::query_as("SELECT $1")
.bind(150_i64) .bind(150_i64)
@@ -28,6 +86,153 @@ pub async fn check_postgres() -> Result<(), sqlx::Error> {
assert_eq!(row.0, 150); assert_eq!(row.0, 150);
info!("Received {}", row.0); 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<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_postgres() -> Result<(), sqlx::Error> {
test_connect().await?;
test_insert_select_update().await?;
Ok(()) Ok(())
} }

View File

@@ -1,14 +1,80 @@
use log::info; use log::info;
use std::collections::HashMap;
use sqlx::sqlite::{ use sqlx::sqlite::{
SqliteConnectOptions, SqliteConnectOptions,
SqlitePool,
SqlitePoolOptions, SqlitePoolOptions,
}; };
pub async fn check_sqlite() -> Result<(), sqlx::Error> { static CREATE_TABLE: &str = "
// "sqlite://test.db" 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<SqlitePool, sqlx::Error> {
// "sqlite://demo.db"
let options = SqliteConnectOptions::new() let options = SqliteConnectOptions::new()
.filename("test.db") .filename("demo.db")
.create_if_missing(true); .create_if_missing(true);
let pool = SqlitePoolOptions::new() let pool = SqlitePoolOptions::new()
@@ -16,6 +82,12 @@ pub async fn check_sqlite() -> Result<(), sqlx::Error> {
.connect_with(options) .connect_with(options)
.await?; .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 // Make a simple query to return the given parameter
let row: (i64,) = sqlx::query_as("SELECT $1") let row: (i64,) = sqlx::query_as("SELECT $1")
.bind(150_i64) .bind(150_i64)
@@ -28,3 +100,162 @@ pub async fn check_sqlite() -> Result<(), sqlx::Error> {
Ok(()) 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<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_sqlite() -> Result<(), sqlx::Error> {
test_connect().await?;
test_insert_select_update().await?;
Ok(())
}

View File

@@ -37,19 +37,44 @@ async fn main() -> Result<(), sqlx::Error> {
let ts = vec![ let ts = vec![
tokio::spawn(async move { 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 { 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 { tokio::spawn(async move {
check_ms_sql().await.unwrap(); match check_sqlite().await {
}), Ok(()) => {},
tokio::spawn(async move { Err(e) => {
check_maria_db().await.unwrap(); panic!("SQLITE ERROR '{}'", e.to_string())
}), },
tokio::spawn(async move { };
check_sqlite().await.unwrap();
}), }),
]; ];