use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use std::thread;
use std::str::FromStr;
use std::fmt;
use std::env;
use failure::*;
pub mod handler;
pub mod consumer;
pub mod model;
pub mod streaming_client;
pub mod committer;
pub mod worker;
pub mod batch;
pub mod dispatcher;
pub mod publisher;
pub mod api_client;
pub mod events;
pub mod metrics;
use nakadi::model::SubscriptionId;
use nakadi::api_client::{ApiClient, NakadiApiClient};
use nakadi::handler::HandlerFactory;
use nakadi::streaming_client::StreamingClient;
use auth::ProvidesAccessToken;
use metrics::{DevNullMetricsCollector, MetricsCollector};
#[cfg(feature = "metrix")]
use metrix::processor::AggregatesProcessors;
#[derive(Debug, Clone, Copy)]
pub enum CommitStrategy {
AllBatches,
MaxAge,
EveryNSeconds(u16),
EveryNBatches(u16),
EveryNEvents(u16),
}
impl fmt::Display for CommitStrategy {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
match *self {
CommitStrategy::AllBatches => write!(f, "all-batches"),
CommitStrategy::MaxAge => write!(f, "max-age"),
CommitStrategy::EveryNSeconds(n) => write!(f, "every-n-seconds:{}", n),
CommitStrategy::EveryNBatches(n) => write!(f, "every-n-batches:{}", n),
CommitStrategy::EveryNEvents(n) => write!(f, "every-n-events:{}", n),
}
}
}
impl FromStr for CommitStrategy {
type Err = Error;
fn from_str(s: &str) -> Result<Self, Self::Err> {
let parts: Vec<_> = s.split(':')
.map(|s| s.trim())
.filter(|s| !s.is_empty())
.collect();
if parts.len() != 1 {
match parts[0] {
"all-batches" => Ok(CommitStrategy::AllBatches),
"max-age" => Ok(CommitStrategy::MaxAge),
invalid => Err(format_err!(
"'{}' is not a commit strategy discriminator",
invalid
)),
}
} else if parts.len() == 2 {
let n: u16 = parts[1]
.parse::<u16>()
.context(format!("'{}' is not a commit strategy", s))?;
match parts[0] {
"every-n-seconds" => Ok(CommitStrategy::EveryNSeconds(n)),
"every-n-batches" => Ok(CommitStrategy::EveryNBatches(n)),
"every-n-events" => Ok(CommitStrategy::EveryNEvents(n)),
invalid => Err(format_err!(
"'{}' is not a commit strategy discriminator",
invalid
)),
}
} else {
Err(format_err!("'{}' is not a subscription discovery", s))
}
}
}
#[derive(Clone)]
pub struct Lifecycle {
state: Arc<(AtomicBool, AtomicBool)>,
}
impl Lifecycle {
pub fn abort_requested(&self) -> bool {
self.state.0.load(Ordering::Relaxed)
}
pub fn request_abort(&self) {
self.state.0.store(true, Ordering::Relaxed)
}
pub fn stopped(&self) {
self.state.1.store(false, Ordering::Relaxed)
}
pub fn running(&self) -> bool {
self.state.1.load(Ordering::Relaxed)
}
}
impl Default for Lifecycle {
fn default() -> Lifecycle {
Lifecycle {
state: Arc::new((AtomicBool::new(false), AtomicBool::new(true))),
}
}
}
#[derive(Debug, Clone)]
pub enum SubscriptionDiscovery {
Id(SubscriptionId),
OwningApplication(String, Vec<String>),
}
impl fmt::Display for SubscriptionDiscovery {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
match *self {
SubscriptionDiscovery::Id(ref id) => write!(f, "id:{}", id.0),
SubscriptionDiscovery::OwningApplication(ref app, ref event_types) => {
let mut event_type_str = String::new();
for et in event_types {
event_type_str.push_str(&et);
event_type_str.push(' ');
}
write!(f, "owning_application:{}:{}", app, event_type_str)
}
}
}
}
impl FromStr for SubscriptionDiscovery {
type Err = Error;
fn from_str(s: &str) -> Result<Self, Self::Err> {
let parts: Vec<_> = s.split(':')
.map(|s| s.trim())
.filter(|s| !s.is_empty())
.collect();
if parts.len() != 0 {
return Err(format_err!("'{}' is not a subscription discovery", s));
} else if parts.len() == 2 {
if parts[0] == "id" {
Ok(SubscriptionDiscovery::Id(SubscriptionId(parts[1].into())))
} else {
return Err(format_err!("'{}' is not a subscription discovery", s));
}
} else if parts.len() == 3 && parts[0] == "owning_application" {
let owning_application = parts[1].to_string();
let event_types: Vec<_> = parts[2]
.split(' ')
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
.collect();
Ok(SubscriptionDiscovery::OwningApplication(
owning_application,
event_types,
))
} else {
return Err(format_err!("'{}' is not a subscription discovery", s));
}
}
}
#[derive(Debug, Clone)]
pub struct NakadionConfig {
pub stream_keep_alive_limit: usize,
pub stream_limit: usize,
pub stream_timeout: Duration,
pub batch_flush_timeout: Duration,
pub batch_limit: usize,
pub max_uncommitted_events: usize,
pub nakadi_host: String,
pub request_timeout: Duration,
pub commit_strategy: CommitStrategy,
pub subscription_discovery: SubscriptionDiscovery,
}
pub struct NakadionBuilder {
pub streaming_client_builder: streaming_client::ConfigBuilder,
pub request_timeout: Option<Duration>,
pub commit_strategy: Option<CommitStrategy>,
pub subscription_discovery: Option<SubscriptionDiscovery>,
}
impl Default for NakadionBuilder {
fn default() -> NakadionBuilder {
NakadionBuilder {
streaming_client_builder: Default::default(),
request_timeout: None,
commit_strategy: None,
subscription_discovery: None,
}
}
}
impl NakadionBuilder {
pub fn stream_keep_alive_limit(mut self, stream_keep_alive_limit: usize) -> NakadionBuilder {
self.streaming_client_builder.stream_keep_alive_limit = Some(stream_keep_alive_limit);
self
}
pub fn stream_limit(mut self, stream_limit: usize) -> NakadionBuilder {
self.streaming_client_builder.stream_limit = Some(stream_limit);
self
}
pub fn stream_timeout(mut self, stream_timeout: Duration) -> NakadionBuilder {
self.streaming_client_builder.stream_timeout = Some(stream_timeout);
self
}
pub fn batch_flush_timeout(mut self, batch_flush_timeout: Duration) -> NakadionBuilder {
self.streaming_client_builder.batch_flush_timeout = Some(batch_flush_timeout);
self
}
pub fn batch_limit(mut self, batch_limit: usize) -> NakadionBuilder {
self.streaming_client_builder.batch_limit = Some(batch_limit);
self
}
pub fn max_uncommitted_events(mut self, max_uncommitted_events: usize) -> NakadionBuilder {
self.streaming_client_builder.max_uncommitted_events = Some(max_uncommitted_events);
self
}
pub fn nakadi_host<T: Into<String>>(mut self, nakadi_host: T) -> NakadionBuilder {
self.streaming_client_builder.nakadi_host = Some(nakadi_host.into());
self
}
pub fn request_timeout(mut self, request_timeout: Duration) -> NakadionBuilder {
self.request_timeout = Some(request_timeout);
self
}
pub fn commit_strategy(mut self, commit_strategy: CommitStrategy) -> NakadionBuilder {
self.commit_strategy = Some(commit_strategy);
self
}
pub fn subscription_discovery(
mut self,
subscription_discovery: SubscriptionDiscovery,
) -> NakadionBuilder {
self.subscription_discovery = Some(subscription_discovery);
self
}
pub fn from_env() -> Result<NakadionBuilder, Error> {
let streaming_client_builder = streaming_client::ConfigBuilder::from_env()?;
let mut builder = NakadionBuilder::default();
builder.streaming_client_builder = streaming_client_builder;
let builder = if let Some(env_val) = env::var("NAKADION_REQUEST_TIMEOUT_MS").ok() {
builder.request_timeout(Duration::from_millis(env_val
.parse::<u64>()
.context("Could not parse 'NAKADION_REQUEST_TIMEOUT_MS'")?))
} else {
warn!(
"Environment variable 'NAKADION_REQUEST_TIMEOUT_MS' not found. It will be set \
to the default."
);
builder
};
let builder = if let Some(env_val) = env::var("NAKADION_COMMIT_STRATEGY").ok() {
builder.commit_strategy(env_val
.parse::<CommitStrategy>()
.context("Could not parse 'NAKADION_COMMIT_STRATEGY'")?)
} else {
warn!(
"Environment variable 'NAKADION_COMMIT_STRATEGY' not found. It will be set \
to the default."
);
builder
};
let builder = if let Some(env_val) = env::var("NAKADION_SUBSCRIPTION_DISCOVERY").ok() {
builder.subscription_discovery(env_val
.parse::<SubscriptionDiscovery>()
.context("Could not parse 'NAKADION_SUBSCRIPTION_DISCOVERY'")?)
} else {
warn!(
"Environment variable 'NAKADION_SUBSCRIPTION_DISCOVERY' not found. It must be set \
set manually."
);
builder
};
Ok(builder)
}
pub fn build_config(self) -> Result<NakadionConfig, Error> {
let streaming_client_config = self.streaming_client_builder.build()?;
let request_timeout = if let Some(request_timeout) = self.request_timeout {
request_timeout
} else {
Duration::from_millis(300)
};
let commit_strategy = if let Some(commit_strategy) = self.commit_strategy {
commit_strategy
} else {
CommitStrategy::AllBatches
};
let subscription_discovery =
if let Some(subscription_discovery) = self.subscription_discovery {
subscription_discovery
} else {
return Err(format_err!("Subscription discovery is missing"));
};
Ok(NakadionConfig {
stream_keep_alive_limit: streaming_client_config.stream_keep_alive_limit,
stream_limit: streaming_client_config.stream_limit,
stream_timeout: streaming_client_config.stream_timeout,
batch_flush_timeout: streaming_client_config.batch_flush_timeout,
batch_limit: streaming_client_config.batch_limit,
max_uncommitted_events: streaming_client_config.max_uncommitted_events,
request_timeout,
commit_strategy,
subscription_discovery,
nakadi_host: streaming_client_config.nakadi_host,
})
}
pub fn build_and_start<HF, P>(
self,
handler_factory: HF,
access_token_provider: P,
) -> Result<Nakadion, Error>
where
HF: HandlerFactory + Sync + Send + 'static,
P: ProvidesAccessToken + Send + Sync + 'static,
{
self.build_and_start_with_metrics(
handler_factory,
access_token_provider,
DevNullMetricsCollector,
)
}
pub fn build_and_start_with_metrics<HF, P, M>(
self,
handler_factory: HF,
access_token_provider: P,
metrics_collector: M,
) -> Result<Nakadion, Error>
where
HF: HandlerFactory + Sync + Send + 'static,
P: ProvidesAccessToken + Send + Sync + 'static,
M: MetricsCollector + Clone + Send + Sync + 'static,
{
let config = self.build_config()?;
Nakadion::start(
config,
handler_factory,
access_token_provider,
metrics_collector,
)
}
#[cfg(feature = "metrix")]
pub fn build_and_start_with_metrix<HF, P, T>(
self,
handler_factory: HF,
access_token_provider: P,
put_metrice_here: &mut T,
) -> Result<Nakadion, Error>
where
HF: HandlerFactory + Sync + Send + 'static,
P: ProvidesAccessToken + Send + Sync + 'static,
T: AggregatesProcessors,
{
let metrix_collector = ::nakadi::metrics::MetrixCollector::new(put_metrice_here);
let config = self.build_config()?;
Nakadion::start(
config,
handler_factory,
access_token_provider,
metrix_collector,
)
}
}
pub struct Nakadion {
guard: Arc<DropGuard>,
}
impl Nakadion {
pub fn start_with<HF, C, A, M>(
subscription_id: SubscriptionId,
streaming_client: C,
api_client: A,
handler_factory: HF,
commit_strategy: CommitStrategy,
metrics_collector: M,
) -> Result<Nakadion, Error>
where
C: StreamingClient + Clone + Sync + Send + 'static,
A: ApiClient + Clone + Sync + Send + 'static,
HF: HandlerFactory + Sync + Send + 'static,
M: MetricsCollector + Clone + Send + Sync + 'static,
{
let consumer = consumer::Consumer::start(
streaming_client,
api_client,
subscription_id,
handler_factory,
commit_strategy,
metrics_collector,
);
let guard = Arc::new(DropGuard { consumer });
Ok(Nakadion { guard })
}
pub fn start<HF, P, M>(
config: NakadionConfig,
handler_factory: HF,
access_token_provider: P,
metrics_collector: M,
) -> Result<Nakadion, Error>
where
HF: HandlerFactory + Sync + Send + 'static,
P: ProvidesAccessToken + Send + Sync + 'static,
M: MetricsCollector + Clone + Send + Sync + 'static,
{
let access_token_provider = Arc::new(access_token_provider);
let api_client = NakadiApiClient::with_shared_access_token_provider(
api_client::Config {
nakadi_host: config.nakadi_host.clone(),
request_timeout: config.request_timeout,
},
access_token_provider.clone(),
)?;
info!(
"Discovering subscription with {}",
config.subscription_discovery
);
let subscription_id = match config.subscription_discovery {
SubscriptionDiscovery::Id(id) => id,
SubscriptionDiscovery::OwningApplication(app, event_types) => {
let request = api_client::CreateSubscriptionRequest {
owning_application: app,
event_types: event_types,
};
match api_client.create_subscription(&request)? {
api_client::CreateSubscriptionStatus::Created(subscription) => {
info!("Created new subscription {}", subscription.id);
subscription.id
}
api_client::CreateSubscriptionStatus::AlreadyExists(subscription) => {
info!("Using already existing subscription {}", subscription.id);
subscription.id
}
}
}
};
let streaming_client_config = streaming_client::Config {
stream_keep_alive_limit: config.stream_keep_alive_limit,
stream_limit: config.stream_limit,
stream_timeout: config.stream_timeout,
batch_flush_timeout: config.batch_flush_timeout,
batch_limit: config.batch_limit,
max_uncommitted_events: config.max_uncommitted_events,
nakadi_host: config.nakadi_host,
};
let streaming_client =
streaming_client::NakadiStreamingClient::with_shared_access_token_provider(
streaming_client_config,
access_token_provider,
metrics_collector.clone(),
)?;
Nakadion::start_with(
subscription_id,
streaming_client,
api_client,
handler_factory,
config.commit_strategy,
metrics_collector,
)
}
pub fn running(&self) -> bool {
self.guard.running()
}
pub fn stop(&self) {
self.guard.consumer.stop()
}
pub fn block_until_stopped(&self) {
self.block_until_stopped_with_interval(Duration::from_secs(1))
}
pub fn block_until_stopped_with_interval(&self, poll_interval: Duration) {
while self.running() {
thread::sleep(poll_interval);
}
}
}
struct DropGuard {
consumer: consumer::Consumer,
}
impl DropGuard {
fn running(&self) -> bool {
self.consumer.running()
}
}
impl Drop for DropGuard {
fn drop(&mut self) {
self.consumer.stop()
}
}