Тестирование¶
Структура¶
sharp/src/tests.rs Protocol unit tests
server/src/tests/
├── mod.rs
├── multi_sub_tests.rs Multi-subscriber delivery
├── nack_tests.rs NACK handling и retry limit
└── ordering_tests.rs Ordering, sequencing, backpressure
client/tests/
├── basic_tests.rs Basic pub/sub flows
├── nack_tests.rs Publisher NACK handling
├── retry_tests.rs Publisher retry logic
└── common/ Shared test infrastructure
Запуск¶
cargo test # Все тесты
cargo test -p sharp # Только протокол
cargo test -p server # Только сервер
cargo test -p client # Только клиент
cargo test test_window # По имени
Что покрыто¶
Протокол:
- Round-trip serialize/deserialize для всех Data variants
- Размер минимального пакета (Ping = 30 bytes)
- Sequenced packets (seq=0 vs seq=None на wire)
- Sequence-based ack/nack
Сервер — multi-subscriber:
- Все подписчики на очереди получают сообщение
- Cleanup pending delivery после ack всех подписчиков
- Immediate ack при отсутствии подписчиков
- Server retry при subscriber timeout
Сервер — NACK:
- NACK одного подписчика не триггерит resend другим
- Retry limit (после
max_subscriber_retriesNACK'ов — drop) - NACK триггерит немедленный resend (не через timeout)
Сервер — ordering:
- FIFO ordering для single publisher
- Monotonic sequence numbers
- Ordering сохраняется после NACK + retransmit
- Window backpressure (window_size=4, 10 сообщений — все в порядке)
- Независимые seq для разных подписчиков
Клиент:
- Basic pub/sub flows
- Publisher NACK handling
- Publisher retry logic
Паттерн тестирования¶
let handle = start_server_random_port(ServerConfig {
max_subscriber_retries: 3,
subscriber_retry_timeout: Duration::from_millis(300),
subscriber_window_size: 128,
}).await.unwrap();
let addr = format!("127.0.0.1:{}", handle.addr().port());
tokio::time::sleep(Duration::from_millis(50)).await;
let sub = Subscriber::new(addr.clone(), Connection::subscriber("q".into())).await.unwrap();
let mut rx = sub.subscribe().await;
let pub_ = Publisher::new(addr, Connection::publisher("q".into())).await.unwrap();
let pkt = Packet::new(PacketType::Durable, Data::Message(Message::new(Bytes::from("test"))));
pub_.publish(pkt).await;
let received = tokio::time::timeout(Duration::from_secs(2), rx.recv())
.await.unwrap().unwrap();
sub.ack(&received).await;
handle.shutdown().await;