1use partitionline::{ProduceRecord, Producer, ProducerConfig, TlsConfig};
12
13#[tokio::main]
14async fn main() -> partitionline::Result<()> {
15 let bootstrap = std::env::var("KAFKA_BOOTSTRAP").unwrap_or_else(|_| "127.0.0.1:9092".into());
16 let topic = std::env::var("KAFKA_TOPIC").unwrap_or_else(|_| "partitionline".into());
17 let mut tls = TlsConfig::default();
18 if let Ok(ca_path) = std::env::var("TLS_CA_PEM") {
19 tls = tls.ca_pem(tokio::fs::read(&ca_path).await.map_err(|e| {
20 partitionline::Error::protocol(format!("read TLS_CA_PEM {ca_path}: {e}"))
21 })?);
22 }
23 if let Ok(name) = std::env::var("TLS_SERVER_NAME") {
24 if !name.is_empty() {
25 tls = tls.server_name(name);
26 }
27 }
28 if let (Ok(cert_path), Ok(key_path)) = (
29 std::env::var("TLS_CLIENT_CERT_PEM"),
30 std::env::var("TLS_CLIENT_KEY_PEM"),
31 ) {
32 let cert = tokio::fs::read(&cert_path).await.map_err(|e| {
33 partitionline::Error::protocol(format!("read TLS_CLIENT_CERT_PEM {cert_path}: {e}"))
34 })?;
35 let key = tokio::fs::read(&key_path).await.map_err(|e| {
36 partitionline::Error::protocol(format!("read TLS_CLIENT_KEY_PEM {key_path}: {e}"))
37 })?;
38 tls = tls.client_identity(cert, key);
39 }
40
41 let producer = Producer::new(ProducerConfig::bootstrap([bootstrap]).tls(tls)).await?;
42 let md = producer
43 .send(ProduceRecord::to(topic).value(&b"hello over tls"[..]))
44 .await?;
45 println!("{}-{}@{}", md.topic, md.partition, md.offset);
46 producer.close().await?;
47 Ok(())
48}