#![cfg(feature = "sqlite")]
mod _test_common;
define_test_user_for_range!(StreamUser, "stream_users_test");
define_test_user_for_pool!(StreamPoolUser, "stream_pool_users_test");
async fn test_stream_basic_impl(config: &_test_common::DbConfig) {
let db = _test_common::create_db_connection(config).await.unwrap();
let _ = db
.execute_sql("DROP TABLE IF EXISTS stream_users_test")
.await;
db.create_table::<StreamUser>().execute().await.unwrap();
for i in 0..10 {
let user = StreamUser {
id: i + 1,
name: format!("user{}", i),
age: 18 + i,
};
db.insert(&user).execute().await.unwrap();
}
let mut stream = db
.select::<StreamUser>()
.stream()
.into_iter()
.await
.unwrap();
let mut count = 0;
while let Some(user_result) = stream.next().await {
let user = user_result.unwrap();
println!("Streamed user: {:?}", user);
count += 1;
}
assert_eq!(count, 10);
println!("✓ test_stream_basic passed");
}
async fn test_stream_with_filter_impl(config: &_test_common::DbConfig) {
let db = _test_common::create_db_connection(config).await.unwrap();
let _ = db
.execute_sql("DROP TABLE IF EXISTS stream_users_test")
.await;
db.create_table::<StreamUser>().execute().await.unwrap();
for i in 0..10 {
let user = StreamUser {
id: i + 1,
name: format!("user{}", i),
age: 18 + i,
};
db.insert(&user).execute().await.unwrap();
}
let mut stream = db
.select::<StreamUser>()
.filter(|u| u.age.ge(23))
.stream()
.into_iter()
.await
.unwrap();
let mut count = 0;
while let Some(user_result) = stream.next().await {
let user = user_result.unwrap();
assert!(user.age >= 23);
println!("Streamed user (age >= 23): {:?}", user);
count += 1;
}
assert_eq!(count, 5);
println!("✓ test_stream_with_filter passed");
}
async fn test_stream_with_order_and_range_impl(config: &_test_common::DbConfig) {
let db = _test_common::create_db_connection(config).await.unwrap();
let _ = db
.execute_sql("DROP TABLE IF EXISTS stream_users_test")
.await;
db.create_table::<StreamUser>().execute().await.unwrap();
for i in 0..10 {
let user = StreamUser {
id: i + 1,
name: format!("user{}", i),
age: 18 + i,
};
db.insert(&user).execute().await.unwrap();
}
let mut stream = db
.select::<StreamUser>()
.order_by_desc(|u| u.age)
.range(0..3)
.stream()
.into_iter()
.await
.unwrap();
let mut count = 0;
let mut prev_age = i32::MAX;
while let Some(user_result) = stream.next().await {
let user = user_result.unwrap();
assert!(user.age <= prev_age); prev_age = user.age;
println!("Streamed user (ordered): {:?}", user);
count += 1;
}
assert_eq!(count, 3);
println!("✓ test_stream_with_order_and_range passed");
}
async fn test_stream_empty_result_impl(config: &_test_common::DbConfig) {
let db = _test_common::create_db_connection(config).await.unwrap();
let _ = db
.execute_sql("DROP TABLE IF EXISTS stream_users_test")
.await;
db.create_table::<StreamUser>().execute().await.unwrap();
let mut stream = db
.select::<StreamUser>()
.stream()
.into_iter()
.await
.unwrap();
let mut count = 0;
while let Some(_user_result) = stream.next().await {
count += 1;
}
assert_eq!(count, 0);
println!("✓ test_stream_empty_result passed");
}
async fn test_stream_in_transaction_impl(config: &_test_common::DbConfig) {
let db = _test_common::create_db_connection(config).await.unwrap();
let _ = db
.execute_sql("DROP TABLE IF EXISTS stream_users_test")
.await;
db.create_table::<StreamUser>().execute().await.unwrap();
for i in 0..5 {
let user = StreamUser {
id: i + 1,
name: format!("txn_user{}", i),
age: 20 + i,
};
db.insert(&user).execute().await.unwrap();
}
let txn = db.begin().await.unwrap();
let mut stream = txn
.select::<StreamUser>()
.filter(|u| u.age.ge(22))
.stream()
.into_iter()
.await
.unwrap();
let mut count = 0;
while let Some(user_result) = stream.next().await {
let user = user_result.unwrap();
assert!(user.age >= 22);
println!("Streamed user in txn: {:?}", user);
count += 1;
}
assert_eq!(count, 3);
drop(stream);
txn.commit().await.unwrap();
println!("✓ test_stream_in_transaction passed");
}
test_on_sqlite_only!(test_stream_basic_impl);
test_on_sqlite_only!(test_stream_with_filter_impl);
test_on_sqlite_only!(test_stream_with_order_and_range_impl);
test_on_sqlite_only!(test_stream_empty_result_impl);
test_on_sqlite_only!(test_stream_in_transaction_impl);
async fn test_stream_large_dataset_impl(config: &_test_common::DbConfig) {
let db = _test_common::create_db_connection(config).await.unwrap();
let _ = db
.execute_sql("DROP TABLE IF EXISTS stream_users_test")
.await;
db.create_table::<StreamUser>().execute().await.unwrap();
const DATA_SIZE: i32 = 1000;
for i in 0..DATA_SIZE {
let user = StreamUser {
id: i + 1,
name: format!("user{}", i),
age: 18 + (i % 50),
};
db.insert(&user).execute().await.unwrap();
}
let mut stream = db
.select::<StreamUser>()
.stream()
.into_iter()
.await
.unwrap();
let mut count = 0;
while let Some(user_result) = stream.next().await {
let user = user_result.unwrap();
assert!(user.id >= 1 && user.id <= DATA_SIZE);
count += 1;
}
assert_eq!(count, DATA_SIZE as usize);
println!(
"✓ test_stream_large_dataset passed ({} records streamed)",
count
);
}
async fn test_stream_error_handling_impl(config: &_test_common::DbConfig) {
let db = _test_common::create_db_connection(config).await.unwrap();
let _ = db
.execute_sql("DROP TABLE IF EXISTS stream_users_test")
.await;
db.create_table::<StreamUser>().execute().await.unwrap();
for i in 0..5 {
let user = StreamUser {
id: i + 1,
name: format!("user{}", i),
age: 20 + i,
};
db.insert(&user).execute().await.unwrap();
}
let mut stream = db
.select::<StreamUser>()
.stream()
.into_iter()
.await
.unwrap();
let mut count = 0;
while let Some(user_result) = stream.next().await {
assert!(user_result.is_ok());
count += 1;
}
assert_eq!(count, 5);
println!("✓ test_stream_error_handling passed");
}
async fn test_stream_with_connection_pool_impl(config: &_test_common::DbConfig) {
let pool = if config.0 == ormer::DbType::Sqlite {
ormer::Database::create_pool(config.0, config.1)
.range(0..1)
.build()
.await
.unwrap()
} else {
ormer::Database::create_pool(config.0, config.1)
.range(1..3)
.build()
.await
.unwrap()
};
let pooled_conn = pool.get().await.unwrap();
let _ = pooled_conn
.execute_sql("DROP TABLE IF EXISTS stream_pool_users_test")
.await;
pooled_conn
.create_table::<StreamPoolUser>()
.execute()
.await
.unwrap();
for i in 0..10 {
let user = StreamPoolUser {
id: i + 1,
name: format!("pool_user{}", i),
age: 18 + i,
email: Some(format!("user{}@test.com", i)),
};
pooled_conn.insert(&user).execute().await.unwrap();
}
let mut stream = pooled_conn
.stream::<StreamPoolUser>()
.into_iter()
.await
.unwrap();
let mut count = 0;
while let Some(user_result) = stream.next().await {
let user = user_result.unwrap();
println!("Streamed pool user: {:?}", user);
count += 1;
}
assert_eq!(count, 10);
println!("✓ test_stream_with_connection_pool passed");
}
async fn test_stream_with_limit_impl(config: &_test_common::DbConfig) {
let db = _test_common::create_db_connection(config).await.unwrap();
let _ = db
.execute_sql("DROP TABLE IF EXISTS stream_users_test")
.await;
db.create_table::<StreamUser>().execute().await.unwrap();
for i in 0..10 {
let user = StreamUser {
id: i + 1,
name: format!("user{}", i),
age: 18 + i,
};
db.insert(&user).execute().await.unwrap();
}
let mut stream = db
.select::<StreamUser>()
.range(0..5)
.stream()
.into_iter()
.await
.unwrap();
let mut count = 0;
while let Some(user_result) = stream.next().await {
let _user = user_result.unwrap();
count += 1;
}
assert_eq!(count, 5);
println!("✓ test_stream_with_limit passed");
}
test_on_sqlite_only!(test_stream_large_dataset_impl);
test_on_sqlite_only!(test_stream_error_handling_impl);
test_on_sqlite_only!(test_stream_with_connection_pool_impl);
test_on_sqlite_only!(test_stream_with_limit_impl);