Compare commits
10 Commits
c9ddbc70bd
...
master
| Author | SHA1 | Date | |
|---|---|---|---|
| a297eb285f | |||
| ff2d7b6504 | |||
| f9f090605c | |||
| d9691c0cc9 | |||
| 01d2afb3dd | |||
| 4b859492c1 | |||
| 80916d3a49 | |||
| 9723f350f9 | |||
| e49f0c90ce | |||
| bf1ac7f368 |
4
.gitignore
vendored
4
.gitignore
vendored
@@ -1,4 +1,6 @@
|
||||
/target
|
||||
Cargo.lock
|
||||
|
||||
*.log
|
||||
test.db*
|
||||
|
||||
demo.db*
|
||||
|
||||
@@ -12,6 +12,13 @@ log = { version = "^0.4", features = [
|
||||
] }
|
||||
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 = { version = "0.5.6", features = [
|
||||
# tokio + rustls
|
||||
|
||||
4
Makefile
4
Makefile
@@ -5,6 +5,9 @@ all:
|
||||
|
||||
clean:
|
||||
rm -rf target
|
||||
rm -rf Cargo.lock
|
||||
rm -rf *.log
|
||||
rm -rf demo.db*
|
||||
|
||||
build:
|
||||
cargo build
|
||||
@@ -14,6 +17,7 @@ test:
|
||||
|
||||
up:
|
||||
docker-compose up -d
|
||||
rm -rf demo.db*
|
||||
|
||||
down:
|
||||
docker-compose down
|
||||
|
||||
22
README.md
22
README.md
@@ -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
|
||||
|
||||
@@ -4,11 +4,11 @@ USE demo;
|
||||
-- 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) ,
|
||||
payload TEXT ,
|
||||
metadata TEXT ,
|
||||
aggregate_type VARCHAR(256) NOT NULL,
|
||||
aggregate_id VARCHAR(256) NOT NULL,
|
||||
sequence bigint CHECK (sequence >= 0) ,
|
||||
payload TEXT ,
|
||||
metadata TEXT ,
|
||||
timestamp timestamp DEFAULT (CURRENT_TIMESTAMP),
|
||||
PRIMARY KEY (aggregate_type, aggregate_id, sequence)
|
||||
);
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -4,11 +4,11 @@ USE demo;
|
||||
-- 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,
|
||||
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,
|
||||
payload TEXT NOT NULL,
|
||||
metadata TEXT NOT NULL,
|
||||
timestamp timestamp DEFAULT (CURRENT_TIMESTAMP),
|
||||
PRIMARY KEY (aggregate_type, aggregate_id, sequence)
|
||||
);
|
||||
|
||||
@@ -8,7 +8,7 @@ services:
|
||||
networks:
|
||||
- default
|
||||
ports:
|
||||
- "8084:3306"
|
||||
- "8081:3306"
|
||||
environment:
|
||||
#- "MARIADB_USER=root"
|
||||
- "MARIADB_ROOT_PASSWORD=admin_pass"
|
||||
@@ -24,7 +24,7 @@ services:
|
||||
networks:
|
||||
- default
|
||||
ports:
|
||||
- "8085:1433"
|
||||
- "8082:1433"
|
||||
environment:
|
||||
#- "SA_USER=sa"
|
||||
- "SA_PASSWORD=adm1n_pa%s"
|
||||
@@ -37,7 +37,7 @@ services:
|
||||
networks:
|
||||
- default
|
||||
ports:
|
||||
- "8086:3306"
|
||||
- "8083:3306"
|
||||
environment:
|
||||
#- "MYSQL_USER=root"
|
||||
- "MYSQL_ROOT_PASSWORD=admin_pass"
|
||||
@@ -50,7 +50,7 @@ services:
|
||||
networks:
|
||||
- default
|
||||
ports:
|
||||
- "8087:5432"
|
||||
- "8084:5432"
|
||||
environment:
|
||||
- "POSTGRES_USER=admin"
|
||||
- "POSTGRES_PASSWORD=admin_pass"
|
||||
@@ -61,9 +61,8 @@ services:
|
||||
networks:
|
||||
- default
|
||||
ports:
|
||||
- "8083:8080"
|
||||
- "8080:8080"
|
||||
depends_on:
|
||||
- maria-db
|
||||
- mssql-db
|
||||
- mysql-db
|
||||
- postgres-db
|
||||
|
||||
@@ -1,15 +1,67 @@
|
||||
use log::info;
|
||||
|
||||
use std::collections::HashMap;
|
||||
|
||||
use sqlx::mysql::{
|
||||
MySqlConnectOptions,
|
||||
MySqlPool,
|
||||
MySqlPoolOptions,
|
||||
};
|
||||
|
||||
pub async fn check_maria_db() -> Result<(), sqlx::Error> {
|
||||
// "mysql://demo_user:demo_pass@localhost:8084/demo"
|
||||
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(8084)
|
||||
.port(8081)
|
||||
.database("demo")
|
||||
.username("demo_user")
|
||||
.password("demo_pass");
|
||||
@@ -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<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(())
|
||||
}
|
||||
|
||||
@@ -1,27 +1,212 @@
|
||||
use log::info;
|
||||
|
||||
use std::collections::HashMap;
|
||||
|
||||
use sqlx::mssql::{
|
||||
MssqlConnectOptions,
|
||||
MssqlPool,
|
||||
};
|
||||
|
||||
pub async fn check_ms_sql() -> Result<(), sqlx::Error> {
|
||||
// "mssql://sa:'adm1n_pa%s'@localhost:8085/demo"
|
||||
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 = $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()
|
||||
.host("localhost")
|
||||
.port(8085)
|
||||
.port(8082)
|
||||
.database("demo")
|
||||
.username("demo_user")
|
||||
.password("pas$w0rd");
|
||||
|
||||
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<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(())
|
||||
}
|
||||
|
||||
@@ -1,15 +1,67 @@
|
||||
use log::info;
|
||||
|
||||
use std::collections::HashMap;
|
||||
|
||||
use sqlx::mysql::{
|
||||
MySqlConnectOptions,
|
||||
MySqlPool,
|
||||
MySqlPoolOptions,
|
||||
};
|
||||
|
||||
pub async fn check_mysql() -> Result<(), sqlx::Error> {
|
||||
// "mysql://demo_user:demo_pass@localhost:8086/demo"
|
||||
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:8083/demo"
|
||||
let options = MySqlConnectOptions::new()
|
||||
.host("localhost")
|
||||
.port(8086)
|
||||
.port(8083)
|
||||
.database("demo")
|
||||
.username("demo_user")
|
||||
.password("demo_pass");
|
||||
@@ -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<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(())
|
||||
}
|
||||
|
||||
@@ -1,15 +1,67 @@
|
||||
use log::info;
|
||||
|
||||
use std::collections::HashMap;
|
||||
|
||||
use sqlx::postgres::{
|
||||
PgConnectOptions,
|
||||
PgPool,
|
||||
PgPoolOptions,
|
||||
};
|
||||
|
||||
pub async fn check_postgres() -> Result<(), sqlx::Error> {
|
||||
// "postgres://demo_user:demo_pass@localhost:8087/demo"
|
||||
static INSERT_EVENT: &str = "
|
||||
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()
|
||||
.host("localhost")
|
||||
.port(8087)
|
||||
.port(8084)
|
||||
.database("demo")
|
||||
.username("demo_user")
|
||||
.password("demo_pass");
|
||||
@@ -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<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(())
|
||||
}
|
||||
|
||||
@@ -1,14 +1,80 @@
|
||||
use log::info;
|
||||
|
||||
use std::collections::HashMap;
|
||||
|
||||
use sqlx::sqlite::{
|
||||
SqliteConnectOptions,
|
||||
SqlitePool,
|
||||
SqlitePoolOptions,
|
||||
};
|
||||
|
||||
pub async fn check_sqlite() -> Result<(), sqlx::Error> {
|
||||
// "sqlite://test.db"
|
||||
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<SqlitePool, sqlx::Error> {
|
||||
// "sqlite://demo.db"
|
||||
let options = SqliteConnectOptions::new()
|
||||
.filename("test.db")
|
||||
.filename("demo.db")
|
||||
.create_if_missing(true);
|
||||
|
||||
let pool = SqlitePoolOptions::new()
|
||||
@@ -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<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(())
|
||||
}
|
||||
|
||||
43
src/main.rs
43
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())
|
||||
},
|
||||
};
|
||||
}),
|
||||
];
|
||||
|
||||
|
||||
Reference in New Issue
Block a user