use std::sync::Arc;
use async_trait::async_trait;
use rdkafka::config::ClientConfig;
use rdkafka::consumer::{CommitMode, Consumer, StreamConsumer};
use rdkafka::message::Message;
#[async_trait]
pub trait GovernanceHandler: Send + Sync {
async fn apply_approved_proposal(&self, proposal_id: &str) -> anyhow::Result<()>;
async fn release_rejected_proposal(&self, _proposal_id: &str) -> anyhow::Result<()> {
Ok(())
}
}
pub struct GovernanceConsumerConfig {
pub brokers: String,
pub topic: String,
pub group_id: String,
}
pub fn spawn_governance_consumer(
cfg: GovernanceConsumerConfig,
handler: Arc<dyn GovernanceHandler>,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
loop {
match build_consumer(&cfg) {
Ok(consumer) => {
#[cfg(feature = "observability")]
tracing::info!(topic = %cfg.topic, group = %cfg.group_id, "governance consumer connected");
run(&consumer, &handler).await;
}
Err(e) => {
#[cfg(feature = "observability")]
tracing::warn!(error = %e, "governance consumer connect failed; retrying");
let _ = e;
}
}
tokio::time::sleep(std::time::Duration::from_secs(10)).await;
}
})
}
async fn run(consumer: &StreamConsumer, handler: &Arc<dyn GovernanceHandler>) {
loop {
match consumer.recv().await {
Ok(message) => {
if let Some(payload) = message.payload() {
handle(payload, handler).await;
}
if let Err(e) = consumer.commit_message(&message, CommitMode::Async) {
#[cfg(feature = "observability")]
tracing::warn!(error = %e, "governance consumer commit failed");
let _ = e;
}
}
Err(e) => {
#[cfg(feature = "observability")]
tracing::warn!(error = %e, "governance consumer recv error; reconnecting");
let _ = e;
break;
}
}
}
}
async fn handle(payload: &[u8], handler: &Arc<dyn GovernanceHandler>) {
let event: serde_json::Value = match serde_json::from_slice(payload) {
Ok(event) => event,
Err(e) => {
#[cfg(feature = "observability")]
tracing::warn!(error = %e, "unparseable policy event");
let _ = e;
return;
}
};
let Some(quorum) = event.get("QuorumReached").and_then(|v| v.as_object()) else {
return;
};
let Some(proposal_id) = quorum.get("proposal_id").and_then(|v| v.as_str()) else {
return;
};
let status = quorum.get("status").and_then(|v| v.as_str());
let result = match status {
Some("Approved") => handler.apply_approved_proposal(proposal_id).await,
Some("Rejected") | Some("Expired") => {
handler.release_rejected_proposal(proposal_id).await
}
_ => Ok(()),
};
if let Err(e) = result {
#[cfg(feature = "observability")]
tracing::warn!(proposal_id, ?status, error = %e, "handling QuorumReached failed");
let _ = (e, status);
}
}
fn build_consumer(cfg: &GovernanceConsumerConfig) -> anyhow::Result<StreamConsumer> {
let settings = {
let mut s = crate::kafka::KafkaSettings::from_env();
s.brokers = cfg.brokers.clone();
s
};
let mut client = ClientConfig::new();
for (k, v) in settings.security_settings() {
client.set(k, v);
}
let consumer: StreamConsumer = client
.set("group.id", &cfg.group_id)
.set("enable.auto.commit", "false")
.set("auto.offset.reset", "earliest")
.set("isolation.level", "read_committed")
.create()?;
consumer.subscribe(&[cfg.topic.as_str()])?;
Ok(consumer)
}