road-runner-common 0.16.1

Shared Rust utilities for exchange ecosystem backend services.
Documentation
//! Generic follower for cex-policy's governance stream (`cex_policy.events`,
//! externally-tagged `DomainEvent`s). Every governed action (FeeChange,
//! AuthzChange, AuthRiskRuleset, ...) reacts to the same `QuorumReached`
//! shape; implement [`GovernanceHandler`] once per action and this consumer
//! handles the reconnect loop, event parsing, and commit-after-apply.

use std::sync::Arc;

use async_trait::async_trait;
use rdkafka::config::ClientConfig;
use rdkafka::consumer::{CommitMode, Consumer, StreamConsumer};
use rdkafka::message::Message;

/// Applies (or releases) a single governed action once cex-policy resolves its
/// proposal. Implement this for each governed action a service owns; the
/// `action` the handler cares about is whatever it decides to apply in
/// [`Self::apply_approved_proposal`] — a service can either wire up one
/// handler per action (multiple consumers) or dispatch internally on the
/// proposal's resolved `action` field, matching cex-metadata's FeeChange today.
#[async_trait]
pub trait GovernanceHandler: Send + Sync {
    /// A `QuorumReached { status: Approved }` fired for `proposal_id`. Fetch the
    /// proposal (payload included) and apply the mutation it governs. Must be
    /// idempotent — Kafka delivery is at-least-once.
    async fn apply_approved_proposal(&self, proposal_id: &str) -> anyhow::Result<()>;

    /// A `QuorumReached { status: Rejected | Expired }` fired for `proposal_id`.
    /// Default no-op — only override if the proposer side tracks a
    /// `pending_proposal_id` that needs releasing so a new proposal can be
    /// submitted (cex-metadata's FeeChange does this; most handlers don't).
    async fn release_rejected_proposal(&self, _proposal_id: &str) -> anyhow::Result<()> {
        Ok(())
    }
}

pub struct GovernanceConsumerConfig {
    pub brokers: String,
    /// e.g. `cex_policy.events`.
    pub topic: String,
    /// Consumer group id. Keep this stable across deploys/rewrites — changing
    /// it replays the whole topic from `auto.offset.reset`.
    pub group_id: String,
}

/// Runs forever, reconnecting on any consumer error. Intended to be
/// `tokio::spawn`ed; hold the returned `JoinHandle` if you want to await/abort it.
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;
            }
        }
    }
}

/// cex-policy's `DomainEvent` is externally tagged:
/// `{ "QuorumReached": { proposal_id, status, occurred_at } }`.
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)
}