WIP series lookup test
This commit is contained in:
@@ -775,3 +775,70 @@ async fn psql_play(db: &Database) -> Result<(), Error> {
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_series_by_channel_01() {
|
||||
let fut = async {
|
||||
use crate as dbpg;
|
||||
let backend = "bck-test-00";
|
||||
let channel = "chn-test-00";
|
||||
let series_by_channel_stats = Arc::new(SeriesByChannelStats::new());
|
||||
let pgconf = test_db_conf();
|
||||
{
|
||||
let (pg, _pg_client_jh) = crate::conn::make_pg_client(&pgconf).await?;
|
||||
crate::schema::schema_check(&pg).await.unwrap();
|
||||
pg.execute("delete from series_by_channel where facility = $1", &[&backend])
|
||||
.await?;
|
||||
}
|
||||
// TODO keep join handles and await later
|
||||
let (channel_info_query_tx, jhs, jh) =
|
||||
dbpg::seriesbychannel::start_lookup_workers(1, &pgconf, series_by_channel_stats.clone()).await?;
|
||||
|
||||
let (tx, rx) = async_channel::bounded(1);
|
||||
let item = dbpg::seriesbychannel::ChannelInfoQuery {
|
||||
backend: backend.into(),
|
||||
channel: channel.into(),
|
||||
scalar_type: netpod::ScalarType::U16.to_scylla_i32(),
|
||||
shape_dims: vec![64],
|
||||
tx: Box::pin(tx),
|
||||
};
|
||||
channel_info_query_tx.send(item).await.unwrap();
|
||||
|
||||
tokio::time::sleep(Duration::from_millis(2000)).await;
|
||||
|
||||
let res = rx.recv().await.unwrap();
|
||||
debug!("received A: {res:?}");
|
||||
|
||||
let (tx, rx) = async_channel::bounded(1);
|
||||
let item = dbpg::seriesbychannel::ChannelInfoQuery {
|
||||
backend: backend.into(),
|
||||
channel: channel.into(),
|
||||
scalar_type: netpod::ScalarType::U16.to_scylla_i32(),
|
||||
shape_dims: vec![64],
|
||||
tx: Box::pin(tx),
|
||||
};
|
||||
channel_info_query_tx.send(item).await.unwrap();
|
||||
let res = rx.recv().await.unwrap();
|
||||
debug!("received B: {res:?}");
|
||||
Ok::<_, Error>(())
|
||||
};
|
||||
taskrun::run(fut).unwrap();
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn test_db_conf() -> Database {
|
||||
Database {
|
||||
host: "127.0.0.1".into(),
|
||||
port: 5432,
|
||||
user: "daqbuffer".into(),
|
||||
pass: "daqbuffer".into(),
|
||||
name: "daqbuffer".into(),
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
async fn test_db_conn() -> Result<PgClient, Error> {
|
||||
let db = test_db_conf();
|
||||
let (pg, pg_client_jh) = crate::conn::make_pg_client(&db).await?;
|
||||
Ok(pg)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user