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) => {
let applied = match message.payload() {
Some(payload) => handle(payload, handler).await,
None => true,
};
if !applied {
break;
}
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>) -> bool {
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 true;
}
};
let Some(quorum) = event.get("QuorumReached").and_then(|v| v.as_object()) else {
return true;
};
let Some(proposal_id) = quorum.get("proposal_id").and_then(|v| v.as_str()) else {
return true;
};
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(()),
};
match result {
Ok(()) => true,
Err(e) => {
#[cfg(feature = "observability")]
tracing::error!(
proposal_id, ?status, error = %e,
"applying QuorumReached failed; holding the offset and retrying"
);
let _ = (e, status);
false
}
}
}
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)
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
#[derive(Default)]
struct FakeHandler {
applied: AtomicUsize,
released: AtomicUsize,
fail_apply: bool,
}
#[async_trait]
impl GovernanceHandler for FakeHandler {
async fn apply_approved_proposal(&self, _proposal_id: &str) -> anyhow::Result<()> {
self.applied.fetch_add(1, Ordering::SeqCst);
if self.fail_apply {
anyhow::bail!("cex-policy unreachable");
}
Ok(())
}
async fn release_rejected_proposal(&self, _proposal_id: &str) -> anyhow::Result<()> {
self.released.fetch_add(1, Ordering::SeqCst);
Ok(())
}
}
fn quorum_event(status: &str) -> Vec<u8> {
serde_json::to_vec(&serde_json::json!({
"QuorumReached": {
"proposal_id": "11111111-1111-4111-8111-111111111111",
"status": status,
"occurred_at": "2026-08-08T00:00:00Z"
}
}))
.unwrap()
}
#[tokio::test]
async fn an_applied_proposal_advances_the_offset() {
let handler: Arc<dyn GovernanceHandler> = Arc::new(FakeHandler::default());
assert!(handle(&quorum_event("Approved"), &handler).await);
}
#[tokio::test]
async fn a_handler_asking_to_retry_holds_the_offset() {
let handler: Arc<dyn GovernanceHandler> =
Arc::new(FakeHandler { fail_apply: true, ..Default::default() });
assert!(
!handle(&quorum_event("Approved"), &handler).await,
"an approved decision the handler could not apply must be redelivered, not discarded"
);
}
#[tokio::test]
async fn rejected_and_expired_go_to_the_release_path() {
for status in ["Rejected", "Expired"] {
let fake = Arc::new(FakeHandler::default());
let handler: Arc<dyn GovernanceHandler> = fake.clone();
assert!(handle(&quorum_event(status), &handler).await);
assert_eq!(fake.released.load(Ordering::SeqCst), 1, "{status}");
assert_eq!(fake.applied.load(Ordering::SeqCst), 0, "{status}");
}
}
#[tokio::test]
async fn unconsumable_events_advance_rather_than_wedge_the_stream() {
let fake = Arc::new(FakeHandler::default());
let handler: Arc<dyn GovernanceHandler> = fake.clone();
assert!(handle(b"not json at all", &handler).await);
assert!(handle(br#"{"BootstrapSealed":{"occurred_at":"2026-08-08T00:00:00Z"}}"#, &handler).await);
assert!(handle(br#"{"QuorumReached":{"status":"Approved"}}"#, &handler).await);
assert!(handle(&quorum_event("Pending"), &handler).await);
assert_eq!(fake.applied.load(Ordering::SeqCst), 0);
assert_eq!(fake.released.load(Ordering::SeqCst), 0);
}
}