improve tests
This commit is contained in:
@@ -12,6 +12,11 @@ log = { version = "^0.4", features = [
|
|||||||
] }
|
] }
|
||||||
fern = "^0.5"
|
fern = "^0.5"
|
||||||
|
|
||||||
|
serde = { version = "^1.0.127", features = ["derive"] }
|
||||||
|
serde_json = "^1.0.66"
|
||||||
|
|
||||||
|
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
|
||||||
|
|||||||
@@ -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;
|
||||||
|
|||||||
@@ -1,11 +1,63 @@
|
|||||||
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> {
|
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<MySqlPool, sqlx::Error> {
|
||||||
// "mysql://demo_user:demo_pass@localhost:8084/demo"
|
// "mysql://demo_user:demo_pass@localhost:8084/demo"
|
||||||
let options = MySqlConnectOptions::new()
|
let options = MySqlConnectOptions::new()
|
||||||
.host("localhost")
|
.host("localhost")
|
||||||
@@ -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(())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,11 +1,62 @@
|
|||||||
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> {
|
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<MssqlPool, sqlx::Error> {
|
||||||
// "mssql://sa:'adm1n_pa%s'@localhost:8085/demo"
|
// "mssql://sa:'adm1n_pa%s'@localhost:8085/demo"
|
||||||
let options = MssqlConnectOptions::new()
|
let options = MssqlConnectOptions::new()
|
||||||
.host("localhost")
|
.host("localhost")
|
||||||
@@ -16,12 +67,146 @@ pub async fn check_ms_sql() -> Result<(), sqlx::Error> {
|
|||||||
|
|
||||||
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(())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,11 +1,63 @@
|
|||||||
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> {
|
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<MySqlPool, sqlx::Error> {
|
||||||
// "mysql://demo_user:demo_pass@localhost:8086/demo"
|
// "mysql://demo_user:demo_pass@localhost:8086/demo"
|
||||||
let options = MySqlConnectOptions::new()
|
let options = MySqlConnectOptions::new()
|
||||||
.host("localhost")
|
.host("localhost")
|
||||||
@@ -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(())
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,11 +1,63 @@
|
|||||||
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> {
|
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<PgPool, sqlx::Error> {
|
||||||
// "postgres://demo_user:demo_pass@localhost:8087/demo"
|
// "postgres://demo_user:demo_pass@localhost:8087/demo"
|
||||||
let options = PgConnectOptions::new()
|
let options = PgConnectOptions::new()
|
||||||
.host("localhost")
|
.host("localhost")
|
||||||
@@ -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(&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
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn check_postgres() -> Result<(), sqlx::Error> {
|
||||||
|
test_connect().await?;
|
||||||
|
test_insert_select_update().await?;
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user