use std::convert::Infallible;
use std::process::ExitCode;
use opentelemetry::KeyValue;
use opentelemetry::global;
use opentelemetry::metrics::Counter;
use ruststream::memory::MemoryBroker;
use ruststream::otel::Otel;
use ruststream::runtime::cli::run_main;
use ruststream::runtime::{App, AppInfo, HandlerResult, RustStream, State};
use ruststream::{FromRef, subscriber};
use serde::Deserialize;
#[derive(Debug, Deserialize)]
struct Order {
#[allow(dead_code)] id: u64,
}
#[derive(Clone)]
struct OrderMetrics {
accepted: Counter<u64>,
}
#[derive(Clone, FromRef)]
struct AppState {
metrics: OrderMetrics,
}
#[subscriber("orders")]
async fn accept(order: &Order, State(metrics): State<OrderMetrics>) -> HandlerResult {
metrics.accepted.add(1, &[KeyValue::new("region", "eu")]);
let _ = order;
HandlerResult::Ack
}
fn app(otel: &Otel) -> impl App + use<> {
RustStream::new(AppInfo::new("orders-svc", "0.1.0"))
.layer(otel.consume_layer())
.publish_layer(otel.publish_layer())
.on_startup(async move |()| {
Ok::<_, Infallible>(AppState {
metrics: OrderMetrics {
accepted: global::meter("orders-svc")
.u64_counter("orders_accepted")
.build(),
},
})
})
.with_broker(MemoryBroker::new(), |b| {
b.include(accept);
})
}
fn main() -> ExitCode {
let otel = Otel::builder()
.service_name("orders-svc")
.otlp_endpoint("http://localhost:4317")
.messaging_system("memory")
.attribute("deployment.environment", "dev")
.init()
.expect("otel init failed");
let code = run_main(|| app(&otel));
otel.shutdown().expect("otel shutdown failed");
code
}