use zksync_concurrency::{ctx, scope, sync, testonly::abort_on_panic};
use zksync_consensus_engine::testonly::TestEngine;
use zksync_consensus_roles::validator;
use super::*;
use crate::testonly;
#[tokio::test]
async fn test_loadtest() {
abort_on_panic();
let ctx = &ctx::test_root(&ctx::RealClock);
let rng = &mut ctx.rng();
let mut setup = validator::testonly::Setup::new(rng, 1);
setup.push_blocks_v1(rng, 10);
let mut cfg = testonly::new_configs(rng, &setup, 0)[0].clone();
cfg.gossip.dynamic_inbound_limit = 7;
scope::run!(ctx, |ctx, s| async {
let engine = TestEngine::new(ctx, &setup).await;
s.spawn_bg(engine.runner.run(ctx));
let (node, runner) = testonly::Instance::new(cfg.clone(), engine.manager.clone());
s.spawn_bg(runner.run(ctx));
for b in &setup.blocks {
engine
.manager
.queue_block(ctx, b.clone())
.await
.context("queue_block()")?;
}
let (send, recv) = ctx::channel::bounded(10);
s.spawn_bg(async {
Loadtest {
addr: cfg.public_addr.clone(),
peer: cfg.gossip.key.public(),
genesis: setup.genesis.clone(),
traffic_pattern: TrafficPattern::Random,
output: Some(send),
}
.run(ctx)
.await?;
Ok(())
});
s.spawn(async {
let mut recv = recv;
let mut count = 0;
while count < 100 {
if recv.recv(ctx).await?.is_some() {
count += 1;
}
}
Ok(())
});
s.spawn(async {
let node = node;
let sub = &mut node.net.gossip.inbound.subscribe();
sync::wait_for(ctx, sub, |pool| {
pool.current().len() == cfg.gossip.dynamic_inbound_limit
})
.await?;
Ok(())
});
Ok(())
})
.await
.unwrap();
}