pub mod pull;
pub mod push;
use std::io::ErrorKind;
use std::time::Duration;
use serde::{Deserialize, Serialize};
use serde_json::json;
use time::serde::rfc3339;
use super::response::Response;
use super::stream::ClusterInfo;
use super::Context;
use crate::jetstream::consumer;
use crate::Error;
pub trait IntoConsumerConfig {
fn into_consumer_config(self) -> Config;
}
#[allow(dead_code)]
#[derive(Clone, Debug)]
pub struct Consumer<T: IntoConsumerConfig> {
pub(crate) context: Context,
pub(crate) config: T,
pub(crate) info: Info,
}
impl<T: IntoConsumerConfig> Consumer<T> {
pub fn new(config: T, info: consumer::Info, context: Context) -> Self {
Self {
config,
info,
context,
}
}
}
impl<T: IntoConsumerConfig> Consumer<T> {
pub async fn info(&mut self) -> Result<&consumer::Info, Error> {
let subject = format!("CONSUMER.INFO.{}.{}", self.info.stream_name, self.info.name);
match self.context.request(subject, &json!({})).await? {
Response::Ok::<Info>(info) => {
self.info = info;
Ok(&self.info)
}
Response::Err { error } => Err(Box::new(std::io::Error::new(
ErrorKind::Other,
format!(
"nats: error while getting consumer info: {}, {}, {}",
error.code, error.status, error.description
),
))),
}
}
async fn fetch_info(&self) -> Result<consumer::Info, Error> {
let subject = format!("CONSUMER.INFO.{}.{}", self.info.stream_name, self.info.name);
match self.context.request(subject, &json!({})).await? {
Response::Ok::<Info>(info) => Ok(info),
Response::Err { error } => Err(Box::new(std::io::Error::new(
ErrorKind::Other,
format!(
"nats: error while getting consumer info: {}, {}, {}",
error.code, error.status, error.description
),
))),
}
}
pub fn cached_info(&self) -> &consumer::Info {
&self.info
}
}
pub trait FromConsumer {
fn try_from_consumer_config(config: crate::jetstream::consumer::Config) -> Result<Self, Error>
where
Self: Sized;
}
pub type PullConsumer = Consumer<self::pull::Config>;
pub type PushConsumer = Consumer<self::push::Config>;
pub type OrderedPushConsumer = Consumer<self::push::OrderedConfig>;
#[derive(Debug, Serialize, Deserialize, Clone, PartialEq, Eq)]
pub struct Info {
pub stream_name: String,
pub name: String,
#[serde(with = "rfc3339")]
pub created: time::OffsetDateTime,
pub config: Config,
pub delivered: SequenceInfo,
pub ack_floor: SequenceInfo,
pub num_ack_pending: usize,
pub num_redelivered: usize,
pub num_waiting: usize,
pub num_pending: u64,
pub cluster: ClusterInfo,
#[serde(default)]
pub push_bound: bool,
}
#[derive(Debug, Serialize, Deserialize, Clone, Copy, PartialEq, Eq)]
pub struct SequenceInfo {
#[serde(rename = "consumer_seq")]
pub consumer_sequence: u64,
#[serde(rename = "stream_seq")]
pub stream_sequence: u64,
#[serde(default, with = "rfc3339::option")]
pub last_active: Option<time::OffsetDateTime>,
}
#[derive(Debug, Default, Serialize, Deserialize, Clone, PartialEq, Eq)]
pub struct Config {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub deliver_subject: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub durable_name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub description: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub deliver_group: Option<String>,
#[serde(flatten)]
pub deliver_policy: DeliverPolicy,
pub ack_policy: AckPolicy,
#[serde(default, with = "serde_nanos", skip_serializing_if = "is_default")]
pub ack_wait: Duration,
#[serde(default, skip_serializing_if = "is_default")]
pub max_deliver: i64,
#[serde(default, skip_serializing_if = "is_default")]
pub filter_subject: String,
pub replay_policy: ReplayPolicy,
#[serde(default, skip_serializing_if = "is_default")]
pub rate_limit: u64,
#[serde(default, skip_serializing_if = "is_default")]
pub sample_frequency: u8,
#[serde(default, skip_serializing_if = "is_default")]
pub max_waiting: i64,
#[serde(default, skip_serializing_if = "is_default")]
pub max_ack_pending: i64,
#[serde(default, skip_serializing_if = "is_default")]
pub headers_only: bool,
#[serde(default, skip_serializing_if = "is_default")]
pub flow_control: bool,
#[serde(default, with = "serde_nanos", skip_serializing_if = "is_default")]
pub idle_heartbeat: Duration,
#[serde(default, skip_serializing_if = "is_default")]
pub max_batch: i64,
#[serde(default, with = "serde_nanos", skip_serializing_if = "is_default")]
pub max_expires: Duration,
#[serde(default, with = "serde_nanos", skip_serializing_if = "is_default")]
pub inactive_threshold: Duration,
#[serde(default, skip_serializing_if = "is_default")]
pub num_replicas: usize,
#[serde(default, skip_serializing_if = "is_default")]
pub memory_storage: bool,
}
impl From<&Config> for Config {
fn from(cc: &Config) -> Config {
cc.clone()
}
}
impl From<&str> for Config {
fn from(s: &str) -> Config {
Config {
durable_name: Some(s.to_string()),
..Default::default()
}
}
}
impl IntoConsumerConfig for Config {
fn into_consumer_config(self) -> Config {
self
}
}
impl IntoConsumerConfig for &Config {
fn into_consumer_config(self) -> Config {
self.clone()
}
}
impl FromConsumer for Config {
fn try_from_consumer_config(config: Config) -> Result<Self, Error>
where
Self: Sized,
{
Ok(config)
}
}
#[derive(Debug, Serialize, Deserialize, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
#[serde(tag = "deliver_policy")]
pub enum DeliverPolicy {
#[serde(rename = "all")]
All,
#[serde(rename = "last")]
Last,
#[serde(rename = "new")]
New,
#[serde(rename = "by_start_sequence")]
ByStartSequence {
#[serde(rename = "opt_start_seq")]
start_sequence: u64,
},
#[serde(rename = "by_start_time")]
ByStartTime {
#[serde(rename = "opt_start_time", with = "rfc3339")]
start_time: time::OffsetDateTime,
},
#[serde(rename = "last_per_subject")]
LastPerSubject,
}
impl Default for DeliverPolicy {
fn default() -> DeliverPolicy {
DeliverPolicy::All
}
}
#[derive(Debug, Serialize, Deserialize, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
pub enum AckPolicy {
#[serde(rename = "explicit")]
Explicit = 2,
#[serde(rename = "none")]
None = 0,
#[serde(rename = "all")]
All = 1,
}
impl Default for AckPolicy {
fn default() -> AckPolicy {
AckPolicy::Explicit
}
}
#[derive(Debug, Serialize, Deserialize, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
pub enum ReplayPolicy {
#[serde(rename = "instant")]
Instant = 0,
#[serde(rename = "original")]
Original = 1,
}
impl Default for ReplayPolicy {
fn default() -> ReplayPolicy {
ReplayPolicy::Instant
}
}
fn is_default<T: Default + Eq>(t: &T) -> bool {
t == &T::default()
}