use std::fmt;
use std::str::FromStr;
use serde::{Deserialize, Deserializer, Serialize, Serializer};
use crate::Error;
new_type! {
#[derive(Debug, Clone, Eq, PartialEq, Ord, PartialOrd, Hash, Serialize, Deserialize)]
pub struct PartitionId(String, env="PARTITION_ID");
}
impl PartitionId {
pub fn join_into_cursor(self, offset: CursorOffset) -> Cursor {
Cursor::new(self, offset)
}
}
impl From<Partition> for PartitionId {
fn from(p: Partition) -> Self {
p.into_partition_id()
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Partition {
pub partition: PartitionId,
pub oldest_available_offset: CursorOffset,
pub newest_available_offset: CursorOffset,
pub unconsumed_events: Option<u64>,
}
impl Partition {
pub fn into_partition_id(self) -> PartitionId {
self.partition
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct Cursor {
pub partition: PartitionId,
pub offset: CursorOffset,
}
impl Cursor {
pub fn new<A: Into<PartitionId>, B: Into<CursorOffset>>(partition: A, offset: B) -> Self {
Cursor {
partition: partition.into(),
offset: offset.into(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum CursorOffset {
Begin,
N(String),
}
impl CursorOffset {
pub fn new<T: AsRef<str>>(offset: T) -> Self {
match offset.as_ref() {
"begin" => CursorOffset::Begin,
other => CursorOffset::N(other.to_string()),
}
}
pub fn as_str(&self) -> &str {
match self {
CursorOffset::Begin => "begin",
CursorOffset::N(ref v) => v.as_ref(),
}
}
}
impl fmt::Display for CursorOffset {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
CursorOffset::Begin => write!(f, "begin")?,
CursorOffset::N(ref v) => write!(f, "{}", v)?,
}
Ok(())
}
}
impl FromStr for CursorOffset {
type Err = Error;
fn from_str(s: &str) -> Result<Self, Self::Err> {
Ok(Self::new(s))
}
}
impl Serialize for CursorOffset {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
serializer.serialize_str(self.as_str())
}
}
impl<'de> Deserialize<'de> for CursorOffset {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
let value = String::deserialize(deserializer)?;
Ok(Self::new(value))
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CursorDistanceQuery {
pub initial_cursor: Cursor,
pub final_cursor: Cursor,
}
impl CursorDistanceQuery {
pub fn new<A: Into<Cursor>, B: Into<Cursor>>(initial_cursor: A, final_cursor: B) -> Self {
CursorDistanceQuery {
initial_cursor: initial_cursor.into(),
final_cursor: final_cursor.into(),
}
}
}
#[derive(Debug, Clone, Eq, PartialEq, Serialize, Deserialize)]
pub struct CursorDistanceResult {
#[serde(flatten)]
pub query: CursorDistanceQuery,
pub distance: u64,
}
#[cfg(test)]
mod test {
use super::*;
use crate::partition::CursorOffset;
use serde_json::{self, json};
#[test]
fn cursor_distance_result() {
let json = json!({
"initial_cursor": {
"partition": "the partition",
"offset": "12345",
},
"final_cursor": {
"partition": "another partition",
"offset": "6789",
},
"distance": 12345,
});
let sample = CursorDistanceResult {
query: CursorDistanceQuery {
initial_cursor: Cursor {
partition: PartitionId::new("the partition"),
offset: CursorOffset::new("12345"),
},
final_cursor: Cursor {
partition: PartitionId::new("another partition"),
offset: CursorOffset::new("6789"),
},
},
distance: 12345,
};
assert_eq!(
serde_json::to_value(sample.clone()).unwrap(),
json,
"serialize"
);
assert_eq!(
serde_json::from_value::<CursorDistanceResult>(json).unwrap(),
sample,
"deserialize"
);
}
}