use std::{error::Error, future::Future, sync::Arc};
use reifydb_client::{Frame, FrameColumn, GrpcChange, GrpcClient, GrpcSubscription, SubscriptionConfig, WireFormat};
use reifydb_value::value::duration::Duration;
use tokio::{runtime::Runtime, time::timeout};
use crate::common::{cleanup_server, create_server_instance, start_server_and_get_grpc_port};
mod basic;
mod batch_mixed_op;
mod data_types;
mod filtered;
mod integration;
mod lifecycle;
mod multiple;
mod notifications;
mod reconnect;
mod stress;
pub fn unique_table_name(prefix: &str) -> String {
use std::time::{SystemTime, UNIX_EPOCH};
let timestamp = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos();
format!("{}_{}", prefix, timestamp % 1_000_000_000)
}
pub async fn create_test_table(
client: &GrpcClient,
name: &str,
columns: &[(&str, &str)],
) -> Result<(), Box<dyn Error>> {
let _ = client.admin("create namespace test", None).await;
let cols = columns.iter().map(|(name, typ)| format!("{}: {}", name, typ)).collect::<Vec<_>>().join(", ");
client.admin(&format!("create table test::{} {{ {} }}", name, cols), None).await?;
Ok(())
}
pub async fn recv_with_timeout(sub: &mut GrpcSubscription, timeout_ms: u64) -> Option<GrpcChange> {
match timeout(Duration::from_milliseconds(timeout_ms as i64).unwrap().to_std(), sub.recv()).await {
Ok(result) => result,
Err(_) => None,
}
}
pub async fn recv_multiple_with_timeout(sub: &mut GrpcSubscription, count: usize, timeout_ms: u64) -> Vec<GrpcChange> {
let mut results = Vec::new();
let deadline = tokio::time::Instant::now() + Duration::from_milliseconds(timeout_ms as i64).unwrap().to_std();
while results.len() < count {
let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
if remaining.is_zero() {
break;
}
match timeout(remaining, sub.recv()).await {
Ok(Some(change)) => results.push(change),
Ok(None) => break,
Err(_) => break,
}
}
results
}
pub fn find_column<'a>(frame: &'a Frame, name: &str) -> Option<&'a FrameColumn> {
frame.columns.iter().find(|c| c.name == name)
}
pub struct SubscriptionTestHarness;
impl SubscriptionTestHarness {
pub fn run<F, Fut>(test_fn: F)
where
F: Fn(TestContext) -> Fut + Send + Sync,
Fut: Future<Output = Result<(), Box<dyn Error>>>,
{
let runtime = Arc::new(Runtime::new().unwrap());
let _guard = runtime.enter();
let mut server = create_server_instance(&runtime);
let port = start_server_and_get_grpc_port(&runtime, &mut server).unwrap();
runtime.block_on(async {
let mut client =
GrpcClient::connect(&format!("http://[::1]:{}", port), WireFormat::Rbcf).await.unwrap();
client.authenticate("mysecrettoken");
let ctx = TestContext::new(client);
test_fn(ctx).await.unwrap();
});
cleanup_server(Some(server));
}
}
pub struct TestContext {
pub client: GrpcClient,
table_prefix: String,
}
impl TestContext {
fn new(client: GrpcClient) -> Self {
Self {
client,
table_prefix: unique_table_name("t"),
}
}
pub async fn rql(&self, query: &str) -> Result<(), Box<dyn Error>> {
self.client.command(query, None).await?;
Ok(())
}
pub async fn create_table(&self, name: &str, columns: &str) -> Result<String, Box<dyn Error>> {
let full_name = format!("{}_{}", self.table_prefix, name);
let _ = self.client.admin("create namespace test", None).await;
self.client.admin(&format!("create table test::{} {{ {} }}", full_name, columns), None).await?;
Ok(full_name)
}
pub async fn subscribe(
&self,
table: &str,
config: SubscriptionConfig,
) -> Result<GrpcSubscription, Box<dyn Error>> {
let sub = self.client.subscribe(&format!("from test::{}", table), config).await?;
Ok(sub)
}
pub async fn insert(&self, table: &str, rows: &str) -> Result<(), Box<dyn Error>> {
self.client.command(&format!("INSERT test::{} [{}]", table, rows), None).await?;
Ok(())
}
pub async fn update(&self, table: &str, filter: &str, map: &str) -> Result<(), Box<dyn Error>> {
self.client
.command(&format!("UPDATE test::{} {{ {} }} FILTER {{{}}}", table, map, filter), None)
.await?;
Ok(())
}
pub async fn delete(&self, table: &str, filter: &str) -> Result<(), Box<dyn Error>> {
self.client.command(&format!("DELETE test::{} FILTER {{{}}}", table, filter), None).await?;
Ok(())
}
pub async fn recv(sub: &mut GrpcSubscription) -> Option<GrpcChange> {
recv_with_timeout(sub, 5000).await
}
}