use azure_core::fmt::SafeDebug;
use serde::de::{DeserializeOwned, Error as _};
use serde::{Deserialize, Deserializer};
use std::time::Duration;
fn deserialize_optional_duration_secs<'de, D>(deserializer: D) -> Result<Option<Duration>, D::Error>
where
D: Deserializer<'de>,
{
let secs = Option::<i64>::deserialize(deserializer)?;
Ok(secs.map(|secs| Duration::from_secs(secs.max(0) as u64)))
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Deserialize)]
#[serde(transparent)]
pub struct LogicalSequenceNumber(i64);
impl LogicalSequenceNumber {
pub fn value(&self) -> i64 {
self.0
}
}
impl From<i64> for LogicalSequenceNumber {
fn from(value: i64) -> Self {
Self(value)
}
}
impl From<LogicalSequenceNumber> for i64 {
fn from(value: LogicalSequenceNumber) -> Self {
value.0
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Deserialize)]
#[serde(rename_all = "camelCase")]
#[non_exhaustive]
pub enum ChangeFeedOperationType {
Create,
Replace,
Delete,
#[serde(other)]
Unknown,
}
#[derive(Clone, SafeDebug, Deserialize)]
#[safe(true)]
#[non_exhaustive]
pub struct ChangeFeedMetadata {
#[serde(rename = "operationType", default)]
operation_type: Option<ChangeFeedOperationType>,
#[serde(rename = "lsn", default)]
lsn: Option<LogicalSequenceNumber>,
#[serde(
rename = "crts",
default,
deserialize_with = "deserialize_optional_duration_secs"
)]
conflict_resolution_timestamp: Option<Duration>,
#[serde(rename = "previousImageLSN", default)]
previous_image_lsn: Option<LogicalSequenceNumber>,
#[serde(rename = "timeToLiveExpired", default)]
time_to_live_expired: Option<bool>,
#[serde(rename = "id", default)]
id: Option<String>,
#[serde(rename = "partitionKey", default)]
partition_key: Option<serde_json::Value>,
}
impl ChangeFeedMetadata {
pub fn operation_type(&self) -> Option<ChangeFeedOperationType> {
self.operation_type
}
pub fn lsn(&self) -> Option<LogicalSequenceNumber> {
self.lsn
}
pub fn conflict_resolution_timestamp(&self) -> Option<Duration> {
self.conflict_resolution_timestamp
}
pub fn previous_image_lsn(&self) -> Option<LogicalSequenceNumber> {
self.previous_image_lsn
}
pub fn time_to_live_expired(&self) -> Option<bool> {
self.time_to_live_expired
}
pub fn id(&self) -> Option<&str> {
self.id.as_deref()
}
pub fn partition_key(&self) -> Option<&serde_json::Value> {
self.partition_key.as_ref()
}
}
#[derive(Clone, Debug)]
#[non_exhaustive]
pub struct ChangeFeedItem<T> {
current: Option<T>,
previous: Option<T>,
metadata: Option<ChangeFeedMetadata>,
}
impl<'de, T> Deserialize<'de> for ChangeFeedItem<T>
where
T: DeserializeOwned,
{
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
let value = serde_json::Value::deserialize(deserializer)?;
let is_envelope = value.as_object().is_some_and(|fields| {
!fields.is_empty()
&& fields
.keys()
.all(|key| matches!(key.as_str(), "current" | "previous" | "metadata"))
});
if is_envelope {
#[derive(Deserialize)]
struct Envelope {
current: Option<serde_json::Value>,
previous: Option<serde_json::Value>,
metadata: Option<ChangeFeedMetadata>,
}
fn document<T: DeserializeOwned>(
value: Option<serde_json::Value>,
) -> Result<Option<T>, serde_json::Error> {
match value {
None | Some(serde_json::Value::Null) => Ok(None),
Some(serde_json::Value::Object(map)) if map.is_empty() => Ok(None),
Some(other) => serde_json::from_value(other).map(Some),
}
}
let Envelope {
current,
previous,
metadata,
} = serde_json::from_value(value).map_err(D::Error::custom)?;
Ok(ChangeFeedItem {
current: document(current).map_err(D::Error::custom)?,
previous: document(previous).map_err(D::Error::custom)?,
metadata,
})
} else {
let current = serde_json::from_value(value).map_err(D::Error::custom)?;
Ok(ChangeFeedItem {
current: Some(current),
previous: None,
metadata: None,
})
}
}
}
impl<T> ChangeFeedItem<T> {
pub fn current(&self) -> Option<&T> {
self.current.as_ref()
}
pub fn previous(&self) -> Option<&T> {
self.previous.as_ref()
}
pub fn metadata(&self) -> Option<&ChangeFeedMetadata> {
self.metadata.as_ref()
}
pub fn operation_type(&self) -> Option<ChangeFeedOperationType> {
self.metadata
.as_ref()
.and_then(ChangeFeedMetadata::operation_type)
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde::Deserialize;
use serde_json::json;
#[derive(Clone, Debug, Deserialize, PartialEq)]
struct Doc {
id: String,
#[serde(default)]
value: Option<i64>,
}
#[test]
fn deserializes_create_envelope() {
let envelope = json!({
"current": { "id": "1", "value": 10 },
"metadata": {
"operationType": "create",
"lsn": 100,
"crts": 1720322460
}
});
let item: ChangeFeedItem<Doc> = serde_json::from_value(envelope).unwrap();
assert_eq!(item.operation_type(), Some(ChangeFeedOperationType::Create));
assert_eq!(
item.current(),
Some(&Doc {
id: "1".into(),
value: Some(10)
})
);
assert!(item.previous().is_none());
let metadata = item.metadata().expect("metadata should be present");
assert_eq!(metadata.lsn(), Some(LogicalSequenceNumber::from(100)));
assert_eq!(
metadata.conflict_resolution_timestamp(),
Some(Duration::from_secs(1720322460))
);
assert!(metadata.previous_image_lsn().is_none());
assert!(metadata.time_to_live_expired().is_none());
}
#[test]
fn deserializes_latest_version_envelope_without_metadata() {
let envelope = json!({ "current": { "id": "1", "value": 42 } });
let item: ChangeFeedItem<Doc> = serde_json::from_value(envelope).unwrap();
assert_eq!(
item.current(),
Some(&Doc {
id: "1".into(),
value: Some(42)
})
);
assert!(item.previous().is_none());
assert!(item.metadata().is_none());
assert!(item.operation_type().is_none());
}
#[test]
fn deserializes_flat_non_enveloped_document() {
let flat = json!({ "id": "9", "value": 99 });
let item: ChangeFeedItem<Doc> = serde_json::from_value(flat).unwrap();
assert_eq!(
item.current(),
Some(&Doc {
id: "9".into(),
value: Some(99)
})
);
assert!(item.previous().is_none());
assert!(item.metadata().is_none());
assert!(item.operation_type().is_none());
}
#[test]
fn flat_document_without_optional_fields_still_deserializes() {
let flat = json!({ "id": "10" });
let item: ChangeFeedItem<Doc> = serde_json::from_value(flat).unwrap();
assert_eq!(
item.current(),
Some(&Doc {
id: "10".into(),
value: None
})
);
assert!(item.previous().is_none());
assert!(item.metadata().is_none());
}
#[test]
fn flat_document_with_reserved_field_name_is_not_misread_as_envelope() {
let flat = json!({
"id": "42",
"value": 7,
"metadata": { "author": "bob" }
});
let item: ChangeFeedItem<Doc> = serde_json::from_value(flat).unwrap();
assert_eq!(
item.current(),
Some(&Doc {
id: "42".into(),
value: Some(7)
})
);
assert!(item.previous().is_none());
assert!(item.metadata().is_none());
let flat = json!({ "id": "43", "previous": "unrelated" });
let item: ChangeFeedItem<Doc> = serde_json::from_value(flat).unwrap();
assert_eq!(item.current().map(|d| d.id.as_str()), Some("43"));
assert!(item.previous().is_none());
}
#[test]
fn delete_envelope_with_only_metadata_is_treated_as_envelope() {
let envelope = json!({ "metadata": { "operationType": "delete", "lsn": 400 } });
let item: ChangeFeedItem<Doc> = serde_json::from_value(envelope).unwrap();
assert!(item.current().is_none());
assert!(item.previous().is_none());
assert_eq!(item.operation_type(), Some(ChangeFeedOperationType::Delete));
}
#[test]
fn delete_envelope_without_current_is_treated_as_envelope() {
let envelope = json!({
"previous": { "id": "3", "value": 30 },
"metadata": { "operationType": "delete", "lsn": 300 }
});
let item: ChangeFeedItem<Doc> = serde_json::from_value(envelope).unwrap();
assert!(item.current().is_none());
assert_eq!(item.previous().map(|d| d.id.as_str()), Some("3"));
assert_eq!(item.operation_type(), Some(ChangeFeedOperationType::Delete));
}
#[test]
fn deserializes_replace_envelope_with_previous() {
let envelope = json!({
"current": { "id": "2", "value": 20 },
"previous": { "id": "2", "value": 15 },
"metadata": {
"operationType": "replace",
"lsn": 200,
"crts": 1720322500,
"previousImageLSN": 199
}
});
let item: ChangeFeedItem<Doc> = serde_json::from_value(envelope).unwrap();
assert_eq!(
item.operation_type(),
Some(ChangeFeedOperationType::Replace)
);
assert_eq!(item.current().and_then(|d| d.value), Some(20));
assert_eq!(item.previous().and_then(|d| d.value), Some(15));
assert_eq!(
item.metadata()
.and_then(ChangeFeedMetadata::previous_image_lsn),
Some(LogicalSequenceNumber::from(199))
);
}
#[test]
fn deserializes_delete_envelope_with_previous_and_ttl() {
let envelope = json!({
"previous": { "id": "3", "value": 30 },
"metadata": {
"operationType": "delete",
"lsn": 300,
"timeToLiveExpired": true
}
});
let item: ChangeFeedItem<Doc> = serde_json::from_value(envelope).unwrap();
assert_eq!(item.operation_type(), Some(ChangeFeedOperationType::Delete));
assert!(item.current().is_none());
assert_eq!(item.previous().map(|d| d.id.as_str()), Some("3"));
assert_eq!(
item.metadata()
.and_then(ChangeFeedMetadata::time_to_live_expired),
Some(true)
);
}
#[test]
fn deserializes_delete_envelope_without_previous() {
let envelope = json!({
"metadata": {
"operationType": "delete",
"lsn": 400
}
});
let item: ChangeFeedItem<Doc> = serde_json::from_value(envelope).unwrap();
assert_eq!(item.operation_type(), Some(ChangeFeedOperationType::Delete));
assert!(item.current().is_none());
assert!(item.previous().is_none());
let metadata = item.metadata().expect("metadata should be present");
assert_eq!(metadata.lsn(), Some(LogicalSequenceNumber::from(400)));
assert!(metadata.time_to_live_expired().is_none());
}
#[test]
fn deserializes_delete_envelope_with_empty_current() {
let envelope = json!({
"current": {},
"metadata": {
"operationType": "delete",
"id": "item-1",
"partitionKey": ["tenant-a"]
}
});
let item: ChangeFeedItem<Doc> = serde_json::from_value(envelope).unwrap();
assert_eq!(item.operation_type(), Some(ChangeFeedOperationType::Delete));
assert!(item.current().is_none());
assert!(item.previous().is_none());
let metadata = item.metadata().expect("metadata should be present");
assert_eq!(metadata.id(), Some("item-1"));
assert_eq!(metadata.partition_key(), Some(&json!(["tenant-a"])));
}
#[test]
fn deserializes_latest_version_envelope_with_partial_metadata() {
let envelope = json!({
"current": { "id": "1", "value": 7 },
"metadata": {
"lsn": 100,
"crts": 1720322460
}
});
let item: ChangeFeedItem<Doc> = serde_json::from_value(envelope).unwrap();
assert_eq!(item.current().and_then(|d| d.value), Some(7));
assert!(item.previous().is_none());
let metadata = item.metadata().expect("metadata should be present");
assert!(metadata.operation_type().is_none());
assert!(item.operation_type().is_none());
assert_eq!(metadata.lsn(), Some(LogicalSequenceNumber::from(100)));
assert_eq!(
metadata.conflict_resolution_timestamp(),
Some(Duration::from_secs(1720322460))
);
}
#[test]
fn operation_type_parses_all_variants() {
for (wire, expected) in [
("create", ChangeFeedOperationType::Create),
("replace", ChangeFeedOperationType::Replace),
("delete", ChangeFeedOperationType::Delete),
] {
let parsed: ChangeFeedOperationType = serde_json::from_value(json!(wire)).unwrap();
assert_eq!(parsed, expected);
}
}
#[test]
fn unknown_operation_type_maps_to_unknown_variant() {
let parsed: ChangeFeedOperationType = serde_json::from_value(json!("resurrect")).unwrap();
assert_eq!(parsed, ChangeFeedOperationType::Unknown);
let envelope = json!({
"current": { "id": "1", "value": 1 },
"metadata": { "operationType": "resurrect", "lsn": 500 }
});
let item: ChangeFeedItem<Doc> = serde_json::from_value(envelope).unwrap();
assert_eq!(
item.operation_type(),
Some(ChangeFeedOperationType::Unknown)
);
assert_eq!(item.current().and_then(|d| d.value), Some(1));
}
}