use std::io::{BufReader, BufWriter, Write};
use std::net::{SocketAddr, TcpStream};
use std::sync::mpsc;
use std::thread::{self, JoinHandle};
use std::time::Duration;
use tephra::log::set::{SegmentConfig, SegmentSet};
use tephra::writer::{WriteCoordinator, WriterConfig};
use tempfile::TempDir;
use tephra_client::{
AppendCondition, AsyncClient, Client, ClientError, ErrorCode, Event, Position, Query,
QueryItem, SequencedEvent, SubEvent, Tag, Tags,
};
use tephra_proto::{DEFAULT_MAX_FRAME_LEN, read_frame, tephra as pb, write_frame};
use tephra_server::{Server, ServerConfig, ShutdownHandle};
use tokio_stream::StreamExt as _;
struct TestServer {
addr: SocketAddr,
shutdown: ShutdownHandle,
server_thread: Option<JoinHandle<()>>,
coordinator: Option<WriteCoordinator>,
_dir: TempDir,
}
impl TestServer {
fn start() -> TestServer {
TestServer::start_with(
ServerConfig::default(),
16 * 1024 * 1024,
WriterConfig::default(),
)
}
fn start_with(
server_config: ServerConfig,
segment_size: usize,
writer_config: WriterConfig,
) -> TestServer {
let dir = TempDir::new().unwrap();
let set = SegmentSet::open(dir.path(), SegmentConfig::new(segment_size)).unwrap();
let (coordinator, handle) = WriteCoordinator::start(set, writer_config).unwrap();
let server = Server::bind("127.0.0.1:0", handle, server_config).unwrap();
let addr = server.local_addr();
let shutdown = server.shutdown_handle();
let server_thread = thread::spawn(move || server.run().expect("server run"));
TestServer {
addr,
shutdown,
server_thread: Some(server_thread),
coordinator: Some(coordinator),
_dir: dir,
}
}
fn client(&self) -> Client {
Client::connect(self.addr).unwrap()
}
}
impl Drop for TestServer {
fn drop(&mut self) {
self.shutdown.shutdown();
if let Some(thread) = self.server_thread.take() {
let _ = thread.join();
}
if let Some(coordinator) = self.coordinator.take() {
coordinator.shutdown();
}
}
}
fn ev(ty: &str, tags: &[&str], payload: &[u8]) -> Event {
Event::new(ty, tags, payload).unwrap()
}
fn tag_set(tags: &[&str]) -> Tags {
Tags::new(
tags.iter()
.map(|tag| Tag::new(*tag).unwrap())
.collect::<Vec<_>>(),
)
.unwrap()
}
fn tag_query(tags: &[&str]) -> Query {
Query::item(QueryItem::with_tags(tag_set(tags)))
}
fn tag_condition(tags: &[&str]) -> AppendCondition {
AppendCondition::new(tag_query(tags))
}
fn positions(events: &[SequencedEvent]) -> Vec<u64> {
events.iter().map(|e| e.position().get()).collect()
}
fn fields(sequenced: &SequencedEvent) -> (u64, String, Vec<String>, Vec<u8>) {
let ev = sequenced.event();
(
sequenced.position().get(),
ev.event_type().to_string(),
ev.tags().map(str::to_string).collect(),
ev.payload().to_vec(),
)
}
#[test]
fn append_then_read_round_trips_events() {
let ts = TestServer::start();
let mut client = ts.client();
let range = client
.append(
[ev("Enrolled", &["course:c1", "student:s1"], b"payload-1")],
None,
)
.unwrap();
assert_eq!(range.first.get(), 1);
assert_eq!(range.last.get(), 1);
client
.append([ev("Renamed", &["course:c1"], b"payload-2")], None)
.unwrap();
let (events, watermark) = client.read_all(Query::all(), Position::ZERO, None).unwrap();
assert_eq!(watermark.get(), 2);
assert_eq!(events.len(), 2);
let (pos, ty, tags, payload) = fields(&events[0]);
assert_eq!(pos, 1);
assert_eq!(ty, "Enrolled");
assert_eq!(tags, vec!["course:c1", "student:s1"]);
assert_eq!(payload, b"payload-1");
let (pos, ty, _, payload) = fields(&events[1]);
assert_eq!(pos, 2);
assert_eq!(ty, "Renamed");
assert_eq!(payload, b"payload-2");
}
#[test]
fn tag_query_filters_the_read() {
let ts = TestServer::start();
let mut client = ts.client();
client.append([ev("A", &["course:c1"], b"")], None).unwrap();
client.append([ev("B", &["course:c2"], b"")], None).unwrap();
client.append([ev("C", &["course:c1"], b"")], None).unwrap();
let (events, _) = client
.read_all(tag_query(&["course:c1"]), Position::ZERO, None)
.unwrap();
assert_eq!(positions(&events), vec![1, 3]);
}
#[test]
fn read_after_skips_the_prefix() {
let ts = TestServer::start();
let mut client = ts.client();
for i in 0..5 {
client
.append([ev("E", &[&format!("k:{i}")], b"")], None)
.unwrap();
}
let (events, watermark) = client
.read_all(Query::all(), Position::new(2), None)
.unwrap();
assert_eq!(positions(&events), vec![3, 4, 5]);
assert_eq!(watermark.get(), 5);
}
#[test]
fn empty_read_yields_only_a_watermark() {
let ts = TestServer::start();
let mut client = ts.client();
client.append([ev("E", &["k:1"], b"")], None).unwrap();
let (events, watermark) = client
.read_all(Query::all(), Position::new(1), None)
.unwrap();
assert!(events.is_empty());
assert_eq!(watermark.get(), 1);
}
#[test]
fn read_limit_caps_the_result_and_paginates() {
let ts = TestServer::start();
let mut client = ts.client();
let total = 25u64;
for i in 0..total {
client
.append(
[ev("Enrolled", &["student:s0"], format!("e{i}").as_bytes())],
None,
)
.unwrap();
client
.append([ev("Enrolled", &["student:s1"], b"noise")], None)
.unwrap();
}
let query = || tag_query(&["student:s0"]);
let (all, _) = client.read_all(query(), Position::ZERO, None).unwrap();
assert_eq!(all.len() as u64, total);
let (page, _) = client.read_all(query(), Position::ZERO, Some(10)).unwrap();
assert_eq!(positions(&page), positions(&all)[..10].to_vec());
let (over, _) = client
.read_all(query(), Position::ZERO, Some(total + 100))
.unwrap();
assert_eq!(positions(&over), positions(&all));
let page_size = 7;
let mut after = Position::ZERO;
let mut tiled: Vec<u64> = Vec::new();
loop {
let (chunk, _) = client.read_all(query(), after, Some(page_size)).unwrap();
if chunk.is_empty() {
break;
}
after = chunk.last().unwrap().position();
tiled.extend(positions(&chunk));
}
assert_eq!(tiled, positions(&all));
let (none, watermark) = client.read_all(query(), Position::ZERO, Some(0)).unwrap();
assert!(none.is_empty());
assert_eq!(watermark.get(), 2 * total);
}
#[test]
fn durable_append_conflict_is_reported_and_not_retryable() {
let ts = TestServer::start();
let mut client = ts.client();
client
.append(
[ev("Reserved", &["username:alice"], b"{}")],
Some(tag_condition(&["username:alice"])),
)
.unwrap();
let err = client
.append(
[ev("Reserved", &["username:alice"], b"{}")],
Some(tag_condition(&["username:alice"])),
)
.unwrap_err();
match err {
ClientError::Server {
code,
retryable,
conflict_position,
..
} => {
assert_eq!(code, ErrorCode::Conflict);
assert!(!retryable, "a durable conflict is terminal");
assert_eq!(conflict_position, Some(Position::new(1)));
}
other => panic!("expected a server conflict, got {other:?}"),
}
}
#[test]
fn malformed_wire_event_maps_to_bad_request() {
let ts = TestServer::start();
let mut event = pb::Event::new();
event.set_type(""); event.tags_mut().push("k:1");
let mut append = pb::AppendRequest::new();
append.events_mut().push(event);
let mut request = pb::Request::new();
request.set_request_id(1);
request.set_append(append);
let response = send_raw_request(ts.addr, &request);
match response.kind() {
pb::response::KindOneof::Error(err) => assert_eq!(err.code(), pb::ErrorCode::BadRequest),
other => panic!("expected a bad-request error, got {other:?}"),
}
}
#[test]
fn empty_append_maps_to_empty() {
let ts = TestServer::start();
let mut client = ts.client();
let err = client.append(Vec::new(), None).unwrap_err();
match err {
ClientError::Server { code, .. } => assert_eq!(code, ErrorCode::Empty),
other => panic!("expected empty, got {other:?}"),
}
}
#[test]
fn large_streamed_read_returns_every_event_in_order() {
let server_config = ServerConfig {
read_batch_events: 7,
read_batch_bytes: 64,
..ServerConfig::default()
};
let writer_config = WriterConfig {
queue_capacity: 64,
max_batch_records: 64,
max_batch_bytes: 256,
..WriterConfig::default()
};
let ts = TestServer::start_with(server_config, 512, writer_config);
let mut client = ts.client();
let total = 200u64;
for i in 0..total {
client
.append([ev("E", &[&format!("n:{i}")], b"x")], None)
.unwrap();
}
let (events, watermark) = client.read_all(Query::all(), Position::ZERO, None).unwrap();
assert_eq!(watermark.get(), total);
let expected: Vec<u64> = (1..=total).collect();
assert_eq!(positions(&events), expected);
}
#[test]
fn streaming_read_iterator_yields_incrementally() {
let ts = TestServer::start();
let mut client = ts.client();
for i in 0..10 {
client
.append([ev("E", &[&format!("k:{i}")], b"")], None)
.unwrap();
}
let mut stream = client.read(Query::all(), Position::ZERO, None).unwrap();
let mut count = 0;
for item in stream.by_ref() {
item.unwrap();
count += 1;
}
assert_eq!(count, 10);
assert_eq!(stream.watermark(), Some(Position::new(10)));
}
#[test]
fn dropping_a_read_early_keeps_the_connection_usable() {
let server_config = ServerConfig {
read_batch_events: 1,
read_batch_bytes: 1,
..ServerConfig::default()
};
let ts = TestServer::start_with(server_config, 16 * 1024 * 1024, WriterConfig::default());
let mut client = ts.client();
for i in 0..20 {
client
.append([ev("E", &[&format!("k:{i}")], b"")], None)
.unwrap();
}
{
let mut stream = client.read(Query::all(), Position::ZERO, None).unwrap();
let first = stream.next().unwrap().unwrap();
assert_eq!(first.position().get(), 1);
}
let (events, watermark) = client.read_all(Query::all(), Position::ZERO, None).unwrap();
assert_eq!(watermark.get(), 20);
assert_eq!(positions(&events), (1..=20).collect::<Vec<_>>());
}
#[test]
fn oversized_request_gets_a_too_large_error_not_a_disconnect() {
let server_config = ServerConfig {
max_frame_len: 64,
..ServerConfig::default()
};
let ts = TestServer::start_with(server_config, 16 * 1024 * 1024, WriterConfig::default());
let mut client = ts.client();
let big = vec![b'x'; 256];
let err = client.append([ev("E", &["k:1"], &big)], None).unwrap_err();
match err {
ClientError::Server { code, .. } => assert_eq!(code, ErrorCode::TooLarge),
other => panic!("expected a TooLarge server error, got {other:?}"),
}
}
#[test]
fn concurrent_clients_stay_consistent() {
let ts = TestServer::start();
let threads = 4u64;
let per_thread = 50u64;
let handles: Vec<_> = (0..threads)
.map(|t| {
let addr = ts.addr;
thread::spawn(move || {
let mut client = Client::connect(addr).unwrap();
for i in 0..per_thread {
client
.append(
[ev("Appended", &[&format!("t:{t}"), &format!("i:{i}")], b"")],
None,
)
.unwrap();
}
})
})
.collect();
for handle in handles {
handle.join().unwrap();
}
let mut client = ts.client();
let (events, watermark) = client.read_all(Query::all(), Position::ZERO, None).unwrap();
let total = threads * per_thread;
assert_eq!(watermark.get(), total);
assert_eq!(events.len() as u64, total);
let mut sorted = positions(&events);
sorted.sort_unstable();
assert_eq!(sorted, (1..=total).collect::<Vec<_>>());
}
#[test]
fn graceful_shutdown_stops_accepting_and_returns() {
let ts = TestServer::start();
let addr = ts.addr;
let mut client = ts.client();
client.append([ev("E", &["k:1"], b"")], None).unwrap();
ts.shutdown.shutdown();
thread::sleep(std::time::Duration::from_millis(50));
let refused = match Client::connect(addr) {
Err(_) => true,
Ok(mut client) => client.append([ev("E", &["k:2"], b"")], None).is_err(),
};
assert!(refused, "server should refuse work after shutdown");
}
fn spawn_subscriber(
mut client: Client,
query: Query,
after: Position,
items: mpsc::Sender<Result<SubEvent, String>>,
cancel: mpsc::Sender<tephra_client::SubscribeCancel>,
) -> JoinHandle<()> {
thread::spawn(move || {
let (stream, canceller) = client.subscribe(query, after).unwrap();
cancel.send(canceller).unwrap();
for item in stream {
match item {
Ok(event) => {
if items.send(Ok(event)).is_err() {
break;
}
}
Err(err) => {
let _ = items.send(Err(err.to_string()));
break;
}
}
}
})
}
#[test]
fn subscribe_streams_catch_up_then_live() {
let ts = TestServer::start();
let mut appender = ts.client();
appender.append([ev("E", &["k:1"], b"a")], None).unwrap();
appender.append([ev("E", &["k:1"], b"b")], None).unwrap();
let (item_tx, item_rx) = mpsc::channel();
let (cancel_tx, cancel_rx) = mpsc::channel();
let subscriber = spawn_subscriber(
ts.client(),
Query::all(),
Position::ZERO,
item_tx,
cancel_tx,
);
let cancel = cancel_rx.recv().unwrap();
let mut positions = Vec::new();
let mut caught_up = Vec::new();
while positions.len() < 2 {
match item_rx.recv().unwrap() {
Ok(SubEvent::Event(ev)) => positions.push(ev.position().get()),
Ok(SubEvent::CaughtUp(w)) => caught_up.push(w),
Err(err) => panic!("subscription error: {err}"),
}
}
appender.append([ev("E", &["k:1"], b"c")], None).unwrap();
appender.append([ev("E", &["k:1"], b"d")], None).unwrap();
while positions.len() < 4 {
match item_rx.recv().unwrap() {
Ok(SubEvent::Event(ev)) => positions.push(ev.position().get()),
Ok(SubEvent::CaughtUp(w)) => caught_up.push(w),
Err(err) => panic!("subscription error: {err}"),
}
}
assert_eq!(
positions,
vec![1, 2, 3, 4],
"no gap or duplicate across the catch-up/live boundary"
);
assert!(
!caught_up.is_empty(),
"expected at least one caught-up marker at the live edge"
);
assert!(
caught_up.windows(2).all(|w| w[0] <= w[1]),
"caught-up watermarks are non-decreasing"
);
cancel.cancel();
subscriber.join().unwrap();
}
#[test]
fn subscribe_from_mid_position_skips_the_prefix() {
let ts = TestServer::start();
let mut appender = ts.client();
for i in 0..4 {
appender
.append([ev("E", &[&format!("k:{i}")], b"")], None)
.unwrap();
}
let (item_tx, item_rx) = mpsc::channel();
let (cancel_tx, cancel_rx) = mpsc::channel();
let subscriber = spawn_subscriber(
ts.client(),
Query::all(),
Position::new(2),
item_tx,
cancel_tx,
);
let cancel = cancel_rx.recv().unwrap();
let mut positions = Vec::new();
while positions.len() < 2 {
match item_rx.recv().unwrap() {
Ok(SubEvent::Event(ev)) => positions.push(ev.position().get()),
Ok(SubEvent::CaughtUp(_)) => {}
Err(err) => panic!("subscription error: {err}"),
}
}
assert_eq!(positions, vec![3, 4]);
cancel.cancel();
subscriber.join().unwrap();
}
#[test]
fn cancel_ends_a_live_subscription() {
let ts = TestServer::start();
let mut appender = ts.client();
appender.append([ev("E", &["k:1"], b"")], None).unwrap();
let (item_tx, item_rx) = mpsc::channel();
let (cancel_tx, cancel_rx) = mpsc::channel();
let subscriber = spawn_subscriber(
ts.client(),
Query::all(),
Position::ZERO,
item_tx,
cancel_tx,
);
let cancel = cancel_rx.recv().unwrap();
match item_rx.recv().unwrap() {
Ok(SubEvent::Event(ev)) => assert_eq!(ev.position().get(), 1),
other => panic!("expected the first event, got {other:?}"),
}
cancel.cancel();
subscriber.join().unwrap();
}
#[test]
fn idle_subscription_does_not_flood_caught_up_frames() {
let server_config = ServerConfig {
subscribe_wait_tick: Duration::from_millis(20),
..ServerConfig::default()
};
let ts = TestServer::start_with(server_config, 16 * 1024 * 1024, WriterConfig::default());
let (item_tx, item_rx) = mpsc::channel();
let (cancel_tx, cancel_rx) = mpsc::channel();
let subscriber = spawn_subscriber(
ts.client(),
Query::all(),
Position::ZERO,
item_tx,
cancel_tx,
);
let cancel = cancel_rx.recv().unwrap();
match item_rx.recv().unwrap() {
Ok(SubEvent::CaughtUp(w)) => assert_eq!(w.get(), 0),
other => panic!("expected a caught-up marker, got {other:?}"),
}
thread::sleep(Duration::from_millis(200));
match item_rx.try_recv() {
Err(mpsc::TryRecvError::Empty) => {}
other => panic!("idle subscription sent an unexpected extra frame: {other:?}"),
}
let mut appender = ts.client();
appender.append([ev("E", &["k:1"], b"")], None).unwrap();
let mut saw_event = false;
let mut saw_caught_up = false;
for _ in 0..2 {
match item_rx.recv().unwrap() {
Ok(SubEvent::Event(ev)) => {
assert_eq!(ev.position().get(), 1);
saw_event = true;
}
Ok(SubEvent::CaughtUp(w)) => {
assert_eq!(w.get(), 1);
saw_caught_up = true;
}
Err(err) => panic!("subscription error: {err}"),
}
}
assert!(
saw_event && saw_caught_up,
"expected the event and one re-armed caught-up marker"
);
cancel.cancel();
subscriber.join().unwrap();
}
#[test]
fn server_shutdown_ends_an_idle_subscription() {
let ts = TestServer::start();
let mut appender = ts.client();
appender.append([ev("E", &["k:1"], b"")], None).unwrap();
let (item_tx, item_rx) = mpsc::channel();
let (cancel_tx, cancel_rx) = mpsc::channel();
let subscriber = spawn_subscriber(
ts.client(),
Query::all(),
Position::ZERO,
item_tx,
cancel_tx,
);
let _cancel = cancel_rx.recv().unwrap();
match item_rx.recv().unwrap() {
Ok(SubEvent::Event(ev)) => assert_eq!(ev.position().get(), 1),
other => panic!("expected the first event, got {other:?}"),
}
drop(ts);
subscriber.join().unwrap();
}
fn send_raw_request(addr: SocketAddr, request: &pb::Request) -> pb::Response {
let stream = TcpStream::connect(addr).unwrap();
stream.set_nodelay(true).unwrap();
let mut writer = BufWriter::new(stream.try_clone().unwrap());
write_frame(&mut writer, request, DEFAULT_MAX_FRAME_LEN).unwrap();
writer.flush().unwrap();
let mut reader = BufReader::new(stream);
read_frame::<pb::Response, _>(&mut reader, DEFAULT_MAX_FRAME_LEN)
.unwrap()
.expect("server closed without responding")
}
fn append_frame(request_id: u64, ty: &str, tag: &str) -> pb::Request {
let mut event = pb::Event::new();
event.set_type(ty);
event.tags_mut().push(tag.to_string());
event.set_payload(b"p".to_vec());
let mut append = pb::AppendRequest::new();
append.events_mut().push(event);
let mut request = pb::Request::new();
request.set_request_id(request_id);
request.set_append(append);
request
}
fn subscribe_all_frame(request_id: u64, after: u64) -> pb::Request {
let mut query = pb::Query::new();
query.set_all(true);
let mut subscribe = pb::SubscribeRequest::new();
subscribe.set_query(query);
subscribe.set_after(after);
let mut request = pb::Request::new();
request.set_request_id(request_id);
request.set_subscribe(subscribe);
request
}
#[test]
fn pipelined_appends_all_succeed_with_dense_positions() {
let ts = TestServer::start();
let stream = TcpStream::connect(ts.addr).unwrap();
stream.set_nodelay(true).unwrap();
let mut writer = BufWriter::new(stream.try_clone().unwrap());
let mut reader = BufReader::new(stream);
let n = 8u64;
for id in 1..=n {
write_frame(
&mut writer,
&append_frame(id, "E", "k:1"),
DEFAULT_MAX_FRAME_LEN,
)
.unwrap();
}
writer.flush().unwrap();
let mut positions = std::collections::HashMap::new();
for _ in 0..n {
let resp = read_frame::<pb::Response, _>(&mut reader, DEFAULT_MAX_FRAME_LEN)
.unwrap()
.expect("a response per pipelined append");
match resp.kind() {
pb::response::KindOneof::Append(append) => {
positions.insert(resp.request_id(), (append.first(), append.last()));
}
other => panic!("expected an append response, got {other:?}"),
}
}
for id in 1..=n {
assert_eq!(positions.get(&id), Some(&(id, id)), "append {id} position");
}
}
#[test]
fn a_subscription_does_not_block_a_concurrent_append() {
let ts = TestServer::start();
let stream = TcpStream::connect(ts.addr).unwrap();
stream.set_nodelay(true).unwrap();
let mut writer = BufWriter::new(stream.try_clone().unwrap());
let mut reader = BufReader::new(stream);
write_frame(
&mut writer,
&subscribe_all_frame(1, 0),
DEFAULT_MAX_FRAME_LEN,
)
.unwrap();
writer.flush().unwrap();
let first = read_frame::<pb::Response, _>(&mut reader, DEFAULT_MAX_FRAME_LEN)
.unwrap()
.expect("subscription responds");
assert_eq!(first.request_id(), 1);
assert!(matches!(first.kind(), pb::response::KindOneof::CaughtUp(_)));
write_frame(
&mut writer,
&append_frame(2, "E", "k:1"),
DEFAULT_MAX_FRAME_LEN,
)
.unwrap();
writer.flush().unwrap();
let mut saw_append = false;
let mut saw_sub_event = false;
for _ in 0..8 {
if saw_append && saw_sub_event {
break;
}
let resp = read_frame::<pb::Response, _>(&mut reader, DEFAULT_MAX_FRAME_LEN)
.unwrap()
.expect("more frames follow");
match (resp.request_id(), resp.kind()) {
(2, pb::response::KindOneof::Append(append)) => {
assert_eq!((append.first(), append.last()), (1, 1));
saw_append = true;
}
(1, pb::response::KindOneof::ReadEvents(events)) => {
assert_eq!(events.events().len(), 1);
assert_eq!(events.events().get(0).unwrap().position(), 1);
saw_sub_event = true;
}
(1, pb::response::KindOneof::CaughtUp(_)) => {} other => panic!("unexpected frame: {other:?}"),
}
}
assert!(saw_append, "the append was answered while subscribed");
assert!(
saw_sub_event,
"the subscription delivered the appended event"
);
}
#[test]
fn cancel_stops_a_subscription_and_frees_the_connection() {
let ts = TestServer::start();
let stream = TcpStream::connect(ts.addr).unwrap();
stream.set_nodelay(true).unwrap();
let mut writer = BufWriter::new(stream.try_clone().unwrap());
let mut reader = BufReader::new(stream);
write_frame(
&mut writer,
&subscribe_all_frame(1, 0),
DEFAULT_MAX_FRAME_LEN,
)
.unwrap();
writer.flush().unwrap();
let caught = read_frame::<pb::Response, _>(&mut reader, DEFAULT_MAX_FRAME_LEN)
.unwrap()
.expect("subscription responds");
assert!(matches!(
caught.kind(),
pb::response::KindOneof::CaughtUp(_)
));
let mut cancel = pb::CancelRequest::new();
cancel.set_target(1);
let mut cancel_req = pb::Request::new();
cancel_req.set_request_id(99);
cancel_req.set_cancel(cancel);
write_frame(&mut writer, &cancel_req, DEFAULT_MAX_FRAME_LEN).unwrap();
writer.flush().unwrap();
write_frame(
&mut writer,
&append_frame(2, "E", "k:1"),
DEFAULT_MAX_FRAME_LEN,
)
.unwrap();
writer.flush().unwrap();
let mut saw_append = false;
for _ in 0..8 {
let resp = read_frame::<pb::Response, _>(&mut reader, DEFAULT_MAX_FRAME_LEN)
.unwrap()
.expect("append is answered after a cancel");
if resp.request_id() == 2 {
assert!(matches!(resp.kind(), pb::response::KindOneof::Append(_)));
saw_append = true;
break;
}
}
assert!(
saw_append,
"append answered after the subscription was cancelled"
);
}
#[tokio::test]
async fn async_client_appends_and_reads_round_trip() {
let ts = TestServer::start();
let client = AsyncClient::connect(ts.addr).await.unwrap();
client
.append([ev("Enrolled", &["course:c1"], b"one")], None)
.await
.unwrap();
client
.append([ev("Enrolled", &["course:c2"], b"two")], None)
.await
.unwrap();
let (events, watermark) = client
.read_all(Query::all(), Position::ZERO, None)
.await
.unwrap();
assert_eq!(events.len(), 2);
assert_eq!(events[0].position(), Position::new(1));
assert_eq!(events[1].position(), Position::new(2));
assert_eq!(watermark, Position::new(2));
}
#[tokio::test]
async fn async_read_limit_caps_the_result_and_paginates() {
let ts = TestServer::start();
let client = AsyncClient::connect(ts.addr).await.unwrap();
for i in 0..12u64 {
client
.append(
[ev("Enrolled", &["student:s0"], format!("e{i}").as_bytes())],
None,
)
.await
.unwrap();
}
let query = || tag_query(&["student:s0"]);
let (page, _) = client
.read_all(query(), Position::ZERO, Some(5))
.await
.unwrap();
assert_eq!(positions(&page), vec![1, 2, 3, 4, 5]);
let after = page.last().unwrap().position();
let (rest, _) = client.read_all(query(), after, Some(100)).await.unwrap();
assert_eq!(positions(&rest), (6..=12).collect::<Vec<_>>());
}
#[tokio::test]
async fn async_client_pipelines_concurrent_appends() {
let ts = TestServer::start();
let client = AsyncClient::connect(ts.addr).await.unwrap();
let n = 16u64;
let mut set = tokio::task::JoinSet::new();
for i in 0..n {
let client = client.clone();
set.spawn(async move {
client
.append([ev("E", &[&format!("k:{i}")], b"p")], None)
.await
.unwrap()
});
}
let mut firsts = std::collections::BTreeSet::new();
while let Some(joined) = set.join_next().await {
let range = joined.unwrap();
assert_eq!(range.first, range.last, "each append is a single event");
firsts.insert(range.first.get());
}
let expected: std::collections::BTreeSet<u64> = (1..=n).collect();
assert_eq!(firsts, expected);
}
#[tokio::test]
async fn async_client_subscribe_coexists_with_append() {
let ts = TestServer::start();
let client = AsyncClient::connect(ts.addr).await.unwrap();
let mut sub = client.subscribe(Query::all(), Position::ZERO).await;
match sub.next().await.unwrap().unwrap() {
SubEvent::CaughtUp(_) => {}
other => panic!("expected a caught-up marker first, got {other:?}"),
}
let range = client
.append([ev("E", &["k:1"], b"p")], None)
.await
.unwrap();
assert_eq!(range.first, Position::new(1));
loop {
match sub.next().await.unwrap().unwrap() {
SubEvent::Event(event) => {
assert_eq!(event.position(), Position::new(1));
break;
}
SubEvent::CaughtUp(_) => {}
}
}
}
#[tokio::test]
async fn async_client_dropping_a_subscription_cancels_and_frees_the_connection() {
let ts = TestServer::start();
let client = AsyncClient::connect(ts.addr).await.unwrap();
{
let mut sub = client.subscribe(Query::all(), Position::ZERO).await;
match sub.next().await.unwrap().unwrap() {
SubEvent::CaughtUp(_) => {}
other => panic!("expected a caught-up marker, got {other:?}"),
}
}
let range = client
.append([ev("E", &["k:1"], b"p")], None)
.await
.unwrap();
assert_eq!(range.first, Position::new(1));
}
#[test]
fn subscription_budget_rejects_excess_subscriptions() {
let config = ServerConfig {
max_concurrent_subscriptions: 2,
..ServerConfig::default()
};
let ts = TestServer::start_with(config, 16 * 1024 * 1024, WriterConfig::default());
let stream = TcpStream::connect(ts.addr).unwrap();
stream.set_nodelay(true).unwrap();
let mut writer = BufWriter::new(stream.try_clone().unwrap());
let mut reader = BufReader::new(stream);
for id in 1..=3u64 {
write_frame(
&mut writer,
&subscribe_all_frame(id, 0),
DEFAULT_MAX_FRAME_LEN,
)
.unwrap();
}
writer.flush().unwrap();
let mut caught_up = 0;
let mut rejected = 0;
for _ in 0..3 {
let resp = read_frame::<pb::Response, _>(&mut reader, DEFAULT_MAX_FRAME_LEN)
.unwrap()
.expect("three responses");
match resp.kind() {
pb::response::KindOneof::CaughtUp(_) => caught_up += 1,
pb::response::KindOneof::Error(_) => {
assert_eq!(
resp.request_id(),
3,
"the third subscription is the one rejected"
);
rejected += 1;
}
other => panic!("unexpected frame: {other:?}"),
}
}
assert_eq!(caught_up, 2, "two subscriptions were accepted");
assert_eq!(rejected, 1, "the third was rejected");
}