Fetching contributors…
Cannot retrieve contributors at this time
83 lines (76 sloc) 2.61 KB
use metric::{LogLine, Telemetry};
use uuid::Uuid;
/// Supported event encodings.
#[derive(Debug, PartialEq, Serialize, Deserialize, Clone)]
pub enum Encoding {
/// Raw bytes, no encoding.
/// Avro
/// JSON
/// Event: the central cernan datastructure
/// Event is the heart of cernan, the enumeration that cernan works on in all
/// cases. The enumeration fields drive sink / source / filter operations
/// depending on their implementation.
#[derive(PartialEq, Debug, Serialize, Deserialize, Clone)]
pub enum Event {
/// A wrapper for `metric::Telemetry`. See its documentation for more
/// detail.
/// A wrapper for `metric::LogLine`. See its documentation for more
/// detail.
/// A flush pulse signal. The `TimerFlush` keeps a counter of the total
/// flushes made in this cernan's run. See `source::Flush` for the origin of
/// these pulses in cernan operation.
/// Shutdown event which marks the location in the queue after which no
/// more events will appear. It is expected that after receiving this
/// marker the given source will exit cleanly.
/// Raw, encoded bytes.
Raw {
/// Ordering value used by some sinks accepting Raw events.
order_by: u64,
/// Encoding for the included bytes.
encoding: Encoding,
/// Encoded payload.
bytes: Vec<u8>,
/// Connection ID of the source on which this raw event was received
connection_id: Option<Uuid>
impl Event {
/// Determine if an event is a `TimerFlush`.
pub fn is_timer_flush(&self) -> bool {
match *self {
Event::TimerFlush(_) => true,
_ => false,
/// Retrieve the timestamp from an `Event` if such exists. `TimerFlush` has
/// no sensible timestamp -- being itself a mechanism _of_ time, not inside
/// time -- and these `Event`s will always return None.
pub fn timestamp(&self) -> Option<i64> {
match *self {
Event::Telemetry(ref telem) => Some(telem.timestamp),
Event::Log(ref log) => Some(log.time),
Event::TimerFlush(_) | Event::Shutdown | Event::Raw { .. } => None,
impl Event {
/// Create a new `Event::Telemetry` from an existing `metric::Telemetry`.
pub fn new_telemetry(metric: Telemetry) -> Event {
/// Create a new `Event::Log` from an existing `metric::LogLine`.
pub fn new_log(log: LogLine) -> Event {