diff --git a/Cargo.toml b/Cargo.toml index 5ccc743..2888e01 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -4,6 +4,7 @@ version = "0.1.0" authors = ["Julius de Bruijn "] license = "Apache-2.0" readme = "README.md" +edition = "2018" description = "A consumer to send push notifications from Kafka" keywords = ["apns", "fcm", "web-push", "consumer", "kafka"] repository = "https://github.com/xray-tech/xorc-notifications" diff --git a/build.rs b/build.rs index c807f48..1b2b523 100644 --- a/build.rs +++ b/build.rs @@ -1,4 +1,4 @@ -extern crate protoc_rust; +use protoc_rust; fn main() { protoc_rust::run(protoc_rust::Args { diff --git a/examples/send_http_request.rs b/examples/send_http_request.rs index 603d150..6a4ef5f 100644 --- a/examples/send_http_request.rs +++ b/examples/send_http_request.rs @@ -1,10 +1,3 @@ -extern crate common; -extern crate clap; -extern crate rdkafka; -extern crate chrono; -extern crate protobuf; -extern crate futures; - use rdkafka::{ config::ClientConfig, producer::future_producer::{ diff --git a/src/apns2/consumer.rs b/src/apns2/consumer.rs index 4b6ea0f..11609dd 100644 --- a/src/apns2/consumer.rs +++ b/src/apns2/consumer.rs @@ -25,8 +25,8 @@ use common::{ use a2::{client::Endpoint, error::Error}; -use notifier::Notifier; -use producer::ApnsProducer; +use crate::notifier::Notifier; +use crate::producer::ApnsProducer; pub struct ApnsHandler { producer: ApnsProducer, @@ -110,7 +110,7 @@ impl EventHandler for ApnsHandler { &self, key: Option>, event: PushNotification, - ) -> Box + 'static + Send> { + ) -> Box + 'static + Send> { let producer = self.producer.clone(); let timer = RESPONSE_TIMES_HISTOGRAM.start_timer(); @@ -145,7 +145,7 @@ impl EventHandler for ApnsHandler { &self, _: Option>, _: HttpRequest, - ) -> Box + 'static + Send> { + ) -> Box + 'static + Send> { warn!("We don't handle http request events here"); Box::new(ok(())) } diff --git a/src/apns2/main.rs b/src/apns2/main.rs index 580f49a..c997d68 100644 --- a/src/apns2/main.rs +++ b/src/apns2/main.rs @@ -2,18 +2,11 @@ #[macro_use] extern crate slog; #[macro_use] extern crate slog_scope; -extern crate a2; -extern crate common; -extern crate futures; -extern crate heck; -extern crate serde_json; -extern crate tokio_timer; - mod consumer; mod notifier; mod producer; -use consumer::ApnsHandler; +use crate::consumer::ApnsHandler; use std::env; use common::{config::Config, system::System}; diff --git a/src/apns2/producer.rs b/src/apns2/producer.rs index ea5e5ca..3b5a501 100644 --- a/src/apns2/producer.rs +++ b/src/apns2/producer.rs @@ -20,7 +20,7 @@ use common::{ }; use heck::SnakeCase; -use CONFIG; +use crate::CONFIG; pub struct ApnsProducer { producer: ResponseProducer, diff --git a/src/common/config.rs b/src/common/config.rs index 9e99575..28231be 100644 --- a/src/common/config.rs +++ b/src/common/config.rs @@ -1,4 +1,4 @@ -use kafka; +use crate::kafka; use toml; use std::{fs::File, io::prelude::*}; diff --git a/src/common/events/mod.rs b/src/common/events.rs similarity index 100% rename from src/common/events/mod.rs rename to src/common/events.rs diff --git a/src/common/kafka/mod.rs b/src/common/kafka.rs similarity index 100% rename from src/common/kafka/mod.rs rename to src/common/kafka.rs diff --git a/src/common/kafka/request_consumer.rs b/src/common/kafka/request_consumer.rs index 1c11926..c4f8ba5 100644 --- a/src/common/kafka/request_consumer.rs +++ b/src/common/kafka/request_consumer.rs @@ -5,8 +5,8 @@ use rdkafka::{ consumer::{CommitMode, Consumer, stream_consumer::StreamConsumer}, topic_partition_list::{Offset, TopicPartitionList}, }; -use kafka::Config; -use events::{ +use crate::kafka::Config; +use crate::events::{ application::Application, push_notification::PushNotification, http_request::HttpRequest, @@ -34,7 +34,7 @@ pub trait EventHandler { &self, key: Option>, event: PushNotification, - ) -> Box + 'static + Send>; + ) -> Box + 'static + Send>; /// Try to send a http request. If key parameter is set, the response /// will be sent with the same routing key. @@ -42,7 +42,7 @@ pub trait EventHandler { &self, key: Option>, event: HttpRequest, - ) -> Box + 'static + Send>; + ) -> Box + 'static + Send>; /// Handle tenant configuration for connection setup. fn handle_config( @@ -98,7 +98,7 @@ impl RequestConsumer { info!("Starting config processing"); - self.handler(consumer, control, &|msg: BorrowedMessage| { + self.handler(consumer, control, &|msg: BorrowedMessage<'_>| { let convert_key = msg.key().and_then(|key| { String::from_utf8(key.to_vec()).ok() }); @@ -162,7 +162,7 @@ impl RequestConsumer { info!("Starting events processing"); - self.handler(consumer, control, &|msg: BorrowedMessage| { + self.handler(consumer, control, &|msg: BorrowedMessage<'_>| { debug!( "Got message"; "topic" => msg.topic(), @@ -197,7 +197,7 @@ impl RequestConsumer { &self, consumer: StreamConsumer, control: oneshot::Receiver<()>, - process_event: &Fn(BorrowedMessage) -> Result<(), ()> + process_event: &dyn Fn(BorrowedMessage<'_>) -> Result<(), ()> ) -> Result<(), ()> { let mut core = Runtime::new().unwrap(); @@ -219,7 +219,7 @@ impl RequestConsumer { Ok(()) } - fn handle_push(&self, msg: &BorrowedMessage) { + fn handle_push(&self, msg: &BorrowedMessage<'_>) { let event_parsing = msg.payload() .and_then(|payload| parse_from_bytes::(payload).ok()); @@ -242,7 +242,7 @@ impl RequestConsumer { } } - fn handle_http(&self, msg: &BorrowedMessage) { + fn handle_http(&self, msg: &BorrowedMessage<'_>) { let event_parsing = msg.payload() .and_then(|payload| parse_from_bytes::(payload).ok()); @@ -261,7 +261,7 @@ impl RequestConsumer { } } - fn handle_config(&self, msg_id: &str, msg: Option<&BorrowedMessage>) { + fn handle_config(&self, msg_id: &str, msg: Option<&BorrowedMessage<'_>>) { let event = msg .and_then(|msg| msg.payload()) .and_then(|payload| parse_from_bytes::(&payload).ok()); diff --git a/src/common/kafka/response_producer.rs b/src/common/kafka/response_producer.rs index 802cd08..57bbfef 100644 --- a/src/common/kafka/response_producer.rs +++ b/src/common/kafka/response_producer.rs @@ -7,7 +7,7 @@ use rdkafka::{ }, }; -use kafka::Config; +use crate::kafka::Config; use protobuf::Message; use std::sync::Arc; @@ -42,7 +42,7 @@ impl ResponseProducer { pub fn publish( &self, key: Option>, - event: &Message, + event: &dyn Message, ) -> DeliveryFuture { let payload = event.write_to_bytes().unwrap(); diff --git a/src/common/lib.rs b/src/common/lib.rs index 9431ace..1bb6e5c 100644 --- a/src/common/lib.rs +++ b/src/common/lib.rs @@ -5,25 +5,6 @@ #[macro_use] extern crate slog; #[macro_use] extern crate slog_scope; -extern crate a2; -extern crate argparse; -extern crate chan_signal; -extern crate chrono; -extern crate erased_serde; -extern crate futures; -extern crate http; -extern crate hyper; -extern crate protobuf; -extern crate rdkafka; -extern crate serde; -extern crate tokio; -extern crate toml; -extern crate web_push; -extern crate slog_json; -extern crate slog_async; -extern crate slog_term; -extern crate regex; - pub mod config; pub mod events; pub mod kafka; diff --git a/src/common/logger.rs b/src/common/logger.rs index 5403311..f313915 100644 --- a/src/common/logger.rs +++ b/src/common/logger.rs @@ -4,7 +4,7 @@ use slog_term::{TermDecorator, CompactFormat}; use slog_async::Async; use slog_json::Json; -use events::{ +use crate::events::{ push_notification::PushNotification, http_response::HttpResponse, http_request::HttpRequest, @@ -47,7 +47,7 @@ impl Logger { } impl KV for PushNotification { - fn serialize(&self, _record: &Record, serializer: &mut Serializer) -> slog::Result { + fn serialize(&self, _record: &Record<'_>, serializer: &mut dyn Serializer) -> slog::Result { serializer.emit_str("device_token", self.get_device_token())?; serializer.emit_str("universe", self.get_universe())?; serializer.emit_str("correlation_id", self.get_header().get_correlation_id())?; @@ -57,7 +57,7 @@ impl KV for PushNotification { } impl slog::Value for HttpResponse { - fn serialize(&self, _record: &Record, _key: Key, serializer: &mut Serializer) -> slog::Result { + fn serialize(&self, _record: &Record<'_>, _key: Key, serializer: &mut dyn Serializer) -> slog::Result { if self.has_payload() { serializer.emit_str( "status_code", @@ -73,7 +73,7 @@ impl slog::Value for HttpResponse { } impl slog::Value for HttpRequest { - fn serialize(&self, _record: &Record, _key: Key, serializer: &mut Serializer) -> slog::Result { + fn serialize(&self, _record: &Record<'_>, _key: Key, serializer: &mut dyn Serializer) -> slog::Result { serializer.emit_str("correlation_id", self.get_header().get_correlation_id())?; serializer.emit_str("request_type", self.get_request_type().as_ref())?; serializer.emit_str("request_body", self.get_body())?; @@ -114,7 +114,7 @@ impl slog::Value for HttpRequest { } impl KV for Application { - fn serialize(&self, _record: &Record, serializer: &mut Serializer) -> slog::Result { + fn serialize(&self, _record: &Record<'_>, serializer: &mut dyn Serializer) -> slog::Result { serializer.emit_str("app_id", self.get_id())?; if self.has_organization() { diff --git a/src/common/system.rs b/src/common/system.rs index 5b37b02..19ede12 100644 --- a/src/common/system.rs +++ b/src/common/system.rs @@ -1,11 +1,11 @@ use chan_signal::{notify, Signal}; -use config::Config; -use kafka::EventHandler; -use kafka::RequestConsumer; -use metrics::StatisticsServer; +use crate::config::Config; +use crate::kafka::EventHandler; +use crate::kafka::RequestConsumer; +use crate::metrics::StatisticsServer; use std::{thread, thread::JoinHandle, sync::Arc}; use futures::sync::oneshot; -use logger::Logger; +use crate::logger::Logger; use slog_scope; pub struct System; diff --git a/src/fcm/consumer.rs b/src/fcm/consumer.rs index 52688df..c038012 100644 --- a/src/fcm/consumer.rs +++ b/src/fcm/consumer.rs @@ -13,8 +13,8 @@ use common::{ use futures::{Future, future::ok}; use std::sync::RwLock; -use notifier::Notifier; -use producer::FcmProducer; +use crate::notifier::Notifier; +use crate::producer::FcmProducer; pub struct FcmHandler { producer: FcmProducer, @@ -56,7 +56,7 @@ impl EventHandler for FcmHandler { &self, key: Option>, event: PushNotification, - ) -> Box + 'static + Send> { + ) -> Box + 'static + Send> { let timer = RESPONSE_TIMES_HISTOGRAM.start_timer(); CALLBACKS_INFLIGHT.inc(); @@ -86,7 +86,7 @@ impl EventHandler for FcmHandler { &self, _: Option>, _: HttpRequest - ) -> Box + 'static + Send> { + ) -> Box + 'static + Send> { warn!("We don't handle http request events here"); Box::new(ok(())) } diff --git a/src/fcm/main.rs b/src/fcm/main.rs index a2badec..7718612 100644 --- a/src/fcm/main.rs +++ b/src/fcm/main.rs @@ -2,17 +2,12 @@ #[macro_use] extern crate slog; #[macro_use] extern crate slog_scope; -extern crate common; -extern crate fcm; -extern crate futures; - mod consumer; mod notifier; mod producer; use common::{config::Config, system::System}; - -use consumer::FcmHandler; +use crate::consumer::FcmHandler; use std::env; lazy_static! { diff --git a/src/fcm/producer.rs b/src/fcm/producer.rs index 47f2ed2..c7226df 100644 --- a/src/fcm/producer.rs +++ b/src/fcm/producer.rs @@ -11,7 +11,7 @@ use common::{ }; use fcm::response::{FcmError, FcmResponse, ErrorReason::*}; -use CONFIG; +use crate::CONFIG; pub struct FcmProducer { producer: ResponseProducer, diff --git a/src/http_requester/consumer.rs b/src/http_requester/consumer.rs index 5ebd31f..b2d4545 100644 --- a/src/http_requester/consumer.rs +++ b/src/http_requester/consumer.rs @@ -9,8 +9,8 @@ use common::{ }; use futures::{Future, future::ok}; -use requester::Requester; -use producer::HttpResponseProducer; +use crate::requester::Requester; +use crate::producer::HttpResponseProducer; pub struct HttpRequestHandler { producer: HttpResponseProducer, @@ -36,7 +36,7 @@ impl EventHandler for HttpRequestHandler { &self, _: Option>, _: PushNotification, - ) -> Box + 'static + Send> { + ) -> Box + 'static + Send> { warn!("We don't handle push notification events here"); Box::new(ok(())) } @@ -45,7 +45,7 @@ impl EventHandler for HttpRequestHandler { &self, key: Option>, event: HttpRequest, - ) -> Box + 'static + Send> { + ) -> Box + 'static + Send> { let producer = self.producer.clone(); let timer = RESPONSE_TIMES_HISTOGRAM.start_timer(); diff --git a/src/http_requester/main.rs b/src/http_requester/main.rs index a13f43b..f65884b 100644 --- a/src/http_requester/main.rs +++ b/src/http_requester/main.rs @@ -2,24 +2,12 @@ #[macro_use] extern crate slog; #[macro_use] extern crate slog_scope; -extern crate tokio_timer; -extern crate protobuf; -extern crate common; -extern crate fcm; -extern crate futures; -extern crate hyper; -extern crate hyper_tls; -extern crate http; -extern crate bytes; -extern crate chrono; - mod consumer; mod requester; mod producer; use common::{config::Config, system::System}; - -use consumer::HttpRequestHandler; +use crate::consumer::HttpRequestHandler; use std::env; lazy_static! { diff --git a/src/http_requester/producer.rs b/src/http_requester/producer.rs index b7838f1..44e8d4a 100644 --- a/src/http_requester/producer.rs +++ b/src/http_requester/producer.rs @@ -11,9 +11,9 @@ use common::{ metrics::* }; use std::{collections::HashMap, str}; -use requester::{HttpResult, RequestError}; +use crate::requester::{HttpResult, RequestError}; -use CONFIG; +use crate::CONFIG; pub struct HttpResponseProducer { producer: ResponseProducer, diff --git a/src/web_push/consumer.rs b/src/web_push/consumer.rs index d33ab01..46cb5ba 100644 --- a/src/web_push/consumer.rs +++ b/src/web_push/consumer.rs @@ -12,8 +12,8 @@ use common::{ use futures::{Future, future::ok}; use std::sync::RwLock; -use notifier::Notifier; -use producer::WebPushProducer; +use crate::notifier::Notifier; +use crate::producer::WebPushProducer; struct ApiKey { fcm_api_key: Option, @@ -59,7 +59,7 @@ impl EventHandler for WebPushHandler { &self, key: Option>, event: PushNotification, - ) -> Box + 'static + Send> { + ) -> Box + 'static + Send> { let producer = self.producer.clone(); match self.fcm_api_keys.read().unwrap().get(event.get_universe()) { @@ -90,7 +90,7 @@ impl EventHandler for WebPushHandler { &self, _: Option>, _: HttpRequest - ) -> Box + 'static + Send> { + ) -> Box + 'static + Send> { warn!("We don't handle http request events here"); Box::new(ok(())) } diff --git a/src/web_push/main.rs b/src/web_push/main.rs index 7a32f20..af4f17b 100644 --- a/src/web_push/main.rs +++ b/src/web_push/main.rs @@ -2,19 +2,12 @@ #[macro_use] extern crate slog; #[macro_use] extern crate slog_scope; -extern crate common; -extern crate futures; -extern crate hyper; -extern crate tokio_signal; -extern crate web_push; - mod consumer; mod notifier; mod producer; use common::{config::Config, system::System}; - -use consumer::WebPushHandler; +use crate::consumer::WebPushHandler; use std::env; lazy_static! { diff --git a/src/web_push/producer.rs b/src/web_push/producer.rs index da5496d..dfd651f 100644 --- a/src/web_push/producer.rs +++ b/src/web_push/producer.rs @@ -10,7 +10,7 @@ use common::{ metrics::CALLBACKS_COUNTER }; -use CONFIG; +use crate::CONFIG; use web_push::{*, WebPushError::*};