use serde::{Deserialize, Deserializer, Serialize, Serializer, de::Error as _};
use crate::{BearerToken, StreamId, TokenId, TokenPermissions};
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum Visibility {
#[default]
Private,
Public,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum RequestedRetention {
Seconds(u64),
Infinite,
}
impl Serialize for RequestedRetention {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
match self {
Self::Seconds(seconds) => serializer.serialize_u64(*seconds),
Self::Infinite => serializer.serialize_str("infinite"),
}
}
}
impl<'de> Deserialize<'de> for RequestedRetention {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
#[derive(Deserialize)]
#[serde(untagged)]
enum WireRetention {
Seconds(u64),
Name(String),
}
match WireRetention::deserialize(deserializer)? {
WireRetention::Seconds(seconds) => Ok(Self::Seconds(seconds)),
WireRetention::Name(name) if name == "infinite" => Ok(Self::Infinite),
WireRetention::Name(_) => Err(D::Error::custom(
"retention must be seconds or \"infinite\"",
)),
}
}
}
#[derive(Clone, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
pub struct CreateStreamRequest {
#[serde(default)]
pub visibility: Visibility,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub retention_secs: Option<RequestedRetention>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub issue_tokens: Option<Vec<TokenPermissions>>,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct IssuedStreamToken {
pub token_id: TokenId,
pub permissions: TokenPermissions,
#[serde(serialize_with = "crate::ids::serialize_bearer_token")]
pub token: BearerToken,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct CreateStreamResponse {
pub stream_id: StreamId,
pub visibility: Visibility,
pub retention_secs: u64,
pub tokens: Vec<IssuedStreamToken>,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct IssueTokenRequest {
pub permissions: TokenPermissions,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub expires_at: Option<String>,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct IssueTokenResponse {
pub token_id: TokenId,
pub permissions: TokenPermissions,
#[serde(serialize_with = "crate::ids::serialize_bearer_token")]
pub token: BearerToken,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct RevokeTokenRequest {
pub token_id: TokenId,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum StreamTokenStatus {
Active,
Expired,
Revoked,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct StreamTokenSummary {
pub token_id: TokenId,
pub permissions: TokenPermissions,
pub status: StreamTokenStatus,
pub issued_at: String,
pub expires_at: Option<String>,
pub revoked_at: Option<String>,
pub is_current: bool,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct ListTokensResponse {
pub tokens: Vec<StreamTokenSummary>,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct StreamInfoResponse {
pub stream_id: StreamId,
pub basin: String,
pub visibility: Visibility,
pub state: String,
pub retention_secs: u64,
pub active_token_count: usize,
}
#[derive(Clone, Debug, Default, Eq, PartialEq, Serialize, Deserialize)]
pub struct UpdateStreamRequest {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub visibility: Option<Visibility>,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct StreamTailResponse {
pub stream_id: StreamId,
pub next_s2_seq_num: u64,
pub last_timestamp_ms: Option<u64>,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub struct StreamRangeResponse {
pub stream_id: StreamId,
pub first_s2_seq_num: Option<u64>,
pub first_timestamp_ms: Option<u64>,
pub next_s2_seq_num: u64,
pub last_timestamp_ms: Option<u64>,
}
#[cfg(test)]
mod tests {
use serde_json::json;
use super::*;
#[test]
fn omits_absent_create_stream_options() {
let request = CreateStreamRequest::default();
assert_eq!(
serde_json::to_value(request).expect("serialize create request"),
json!({ "visibility": "private" })
);
}
#[test]
fn serializes_finite_and_infinite_retention_requests() {
for (retention, expected) in [
(RequestedRetention::Seconds(604_800), json!(604_800)),
(RequestedRetention::Infinite, json!("infinite")),
] {
let request = CreateStreamRequest {
retention_secs: Some(retention),
..CreateStreamRequest::default()
};
let value = serde_json::to_value(request).expect("serialize create request");
assert_eq!(value["retention_secs"], expected);
assert_eq!(
serde_json::from_value::<CreateStreamRequest>(value)
.expect("deserialize create request")
.retention_secs,
Some(retention)
);
}
assert!(
serde_json::from_value::<CreateStreamRequest>(json!({
"visibility": "private",
"retention_secs": "forever"
}))
.is_err()
);
}
#[test]
fn serializes_token_mutations_and_omits_absent_stream_update() {
let token = IssueTokenRequest {
permissions: TokenPermissions::read(),
expires_at: None,
};
let token_id: TokenId = "0123456789abcdefghjkmnpq".parse().expect("token id");
assert_eq!(
serde_json::to_value(token).expect("serialize token request"),
json!({ "permissions": "r" })
);
assert_eq!(
serde_json::to_value(RevokeTokenRequest { token_id })
.expect("serialize token revocation request"),
json!({ "token_id": "0123456789abcdefghjkmnpq" })
);
assert_eq!(
serde_json::to_value(UpdateStreamRequest::default()).expect("serialize update request"),
json!({})
);
}
}