Skip to main content

wakeup/
wakeup.rs

1//! Interrupt an in-flight fetch from another task.
2//!
3//! Broker on `KAFKA_BOOTSTRAP`. Topic `KAFKA_TOPIC` (default `partitionline`).
4//! Ctrl-C (or any second task holding [`partitionline::WakeupHandle`]) stops
5//! `fetch` with [`partitionline::Error::Wakeup`].
6
7use partitionline::Consumer;
8
9#[tokio::main]
10async fn main() -> partitionline::Result<()> {
11    let bootstrap = std::env::var("KAFKA_BOOTSTRAP").unwrap_or_else(|_| "127.0.0.1:9092".into());
12    let topic = std::env::var("KAFKA_TOPIC").unwrap_or_else(|_| "partitionline".into());
13    let mut consumer = Consumer::connect(bootstrap).await?;
14    consumer.assign_topic(topic, 0).await?;
15    let wakeup = consumer.wakeup_handle();
16    drop(tokio::spawn(async move {
17        tokio::signal::ctrl_c().await.unwrap_or(());
18        wakeup.wakeup();
19    }));
20    loop {
21        match consumer.fetch().await {
22            Ok(recs) => {
23                for rec in recs {
24                    println!("{}-{}@{}", rec.topic, rec.partition, rec.offset);
25                }
26            }
27            Err(partitionline::Error::Wakeup) => break,
28            Err(e) => return Err(e),
29        }
30    }
31    consumer.close().await?;
32    Ok(())
33}