262 lines
5.5 KiB
Rust
262 lines
5.5 KiB
Rust
use log::info;
|
|
|
|
use std::collections::HashMap;
|
|
|
|
use sqlx::sqlite::{
|
|
SqliteConnectOptions,
|
|
SqlitePool,
|
|
SqlitePoolOptions,
|
|
};
|
|
|
|
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://test.db"
|
|
let options = SqliteConnectOptions::new()
|
|
.filename("test.db")
|
|
.create_if_missing(true);
|
|
|
|
let pool = SqlitePoolOptions::new()
|
|
.max_connections(5)
|
|
.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)
|
|
.fetch_one(&pool)
|
|
.await?;
|
|
|
|
assert_eq!(row.0, 150);
|
|
|
|
info!("Received {}", row.0);
|
|
|
|
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(())
|
|
}
|