oura 2.2.0

The tail of Cardano
Documentation
use std::time::Duration;

use gasket::framework::*;
use kafka::producer::{Producer, Record, RequiredAcks};
use serde::Deserialize;

use crate::framework::*;

pub struct Worker {
    producer: Producer,
    partitioning: PartitionStrategy,
}

#[async_trait::async_trait(?Send)]
impl gasket::framework::Worker<Stage> for Worker {
    async fn bootstrap(stage: &Stage) -> Result<Self, WorkerError> {
        let mut builder = Producer::from_hosts(stage.config.brokers.clone());

        if let Some(timeout) = stage.config.ack_timeout_secs {
            builder = builder
                .with_ack_timeout(Duration::from_secs(timeout))
                .with_required_acks(RequiredAcks::One)
        };

        let producer = builder.create().or_panic()?;

        let partitioning = stage
            .config
            .paritioning
            .clone()
            .unwrap_or(PartitionStrategy::Random);

        Ok(Self {
            producer,
            partitioning,
        })
    }

    async fn schedule(
        &mut self,
        stage: &mut Stage,
    ) -> Result<WorkSchedule<ChainEvent>, WorkerError> {
        let msg = stage.input.recv().await.or_panic()?;
        Ok(WorkSchedule::Unit(msg.payload))
    }

    async fn execute(&mut self, unit: &ChainEvent, stage: &mut Stage) -> Result<(), WorkerError> {
        let point = unit.point().clone();
        let record = unit.record().cloned();

        if record.is_none() {
            return Ok(());
        }

        let payload = serde_json::to_vec(&serde_json::Value::from(record.unwrap())).or_panic()?;

        match self.partitioning {
            PartitionStrategy::ByBlock => {
                let slot = point.slot_or_default().to_be_bytes();
                let kafka_record = Record::from_key_value(&stage.config.topic, &slot[..], payload);
                self.producer.send(&kafka_record)
            }
            PartitionStrategy::Random => {
                let kafka_record = Record::from_value(&stage.config.topic, payload);
                self.producer.send(&kafka_record)
            }
        }
        .or_retry()?;

        stage.ops_count.inc(1);
        stage.latest_block.set(point.slot_or_default() as i64);
        stage.cursor.send(point.clone().into()).await.or_panic()?;

        Ok(())
    }
}

#[derive(Stage)]
#[stage(name = "sink-kafka", unit = "ChainEvent", worker = "Worker")]
pub struct Stage {
    config: Config,

    pub input: MapperInputPort,
    pub cursor: SinkCursorPort,

    #[metric]
    ops_count: gasket::metrics::Counter,

    #[metric]
    latest_block: gasket::metrics::Gauge,
}

#[derive(Debug, Clone, Deserialize)]
pub enum PartitionStrategy {
    ByBlock,
    Random,
}

#[derive(Default, Debug, Deserialize)]
pub struct Config {
    pub brokers: Vec<String>,
    pub topic: String,
    pub ack_timeout_secs: Option<u64>,
    pub paritioning: Option<PartitionStrategy>,
}

impl Config {
    pub fn bootstrapper(self, _ctx: &Context) -> Result<Stage, Error> {
        let stage = Stage {
            config: self,
            ops_count: Default::default(),
            latest_block: Default::default(),
            input: Default::default(),
            cursor: Default::default(),
        };

        Ok(stage)
    }
}