use crate::state::{ChannelId, Seq};
use crate::sync_session::IncrementChannel;
pub fn parse_increment_channel(data: &serde_json::Value) -> Option<IncrementChannel> {
let channel_id = data
.get("id")
.and_then(|v| v.as_str())
.and_then(ChannelId::from_str)?;
for key in ["unreadCount", "unread_count"] {
if let Some(value) = data.get(key) {
if value.as_i64().filter(|count| *count >= 0).is_none() {
return None;
}
}
}
let last_event_seq = Seq(data
.get("lastEventSeq")
.and_then(|v| v.as_u64())
.unwrap_or(0));
let need_sync = data
.get("needSync")
.and_then(|v| v.as_bool())
.unwrap_or(true);
let unread_post_id = ["unread_post_id", "unreadPostId", "unReadPostId"]
.into_iter()
.find_map(|key| data.get(key))
.and_then(|value| value.as_str().map(str::to_owned));
let last_read_seq = match ["last_read_seq", "lastReadSeq"]
.into_iter()
.find_map(|key| data.get(key))
{
None => None,
Some(value) => {
let read_seq = parse_i64(value)?;
if read_seq < 0 {
return None;
}
Some(read_seq)
}
};
let projection_revision = match ["projection_revision", "projectionRevision"]
.into_iter()
.find_map(|key| data.get(key))
{
None => None,
Some(value) => Some(parse_u64(value)?),
};
let raw = bytes::Bytes::from(serde_json::to_vec(data).ok()?);
Some(IncrementChannel {
channel_id,
last_event_seq,
need_sync,
unread_post_id,
last_read_seq,
projection_revision,
raw,
})
}
fn parse_i64(value: &serde_json::Value) -> Option<i64> {
value
.as_i64()
.or_else(|| value.as_u64().and_then(|number| i64::try_from(number).ok()))
.or_else(|| value.as_str().and_then(|number| number.parse::<i64>().ok()))
}
fn parse_u64(value: &serde_json::Value) -> Option<u64> {
value
.as_u64()
.or_else(|| value.as_i64().and_then(|number| u64::try_from(number).ok()))
.or_else(|| value.as_str().and_then(|number| number.parse::<u64>().ok()))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parses_canonical_unread_projection_fields_and_preserves_empty_anchor() {
let data = serde_json::json!({
"id": "bi643wys7fgy9pbai9wjfusnfa",
"lastEventSeq": 151,
"needSync": false,
"unread_post_id": "",
"last_read_seq": 41,
"projection_revision": 7
});
let parsed = parse_increment_channel(&data).expect("valid increment frame");
assert_eq!(parsed.unread_post_id.as_deref(), Some(""));
assert_eq!(parsed.last_read_seq, Some(41));
assert_eq!(parsed.projection_revision, Some(7));
}
#[test]
fn malformed_versioned_integer_is_rejected_instead_of_becoming_legacy() {
let data = serde_json::json!({
"id": "bi643wys7fgy9pbai9wjfusnfa",
"projectionRevision": "not-a-number"
});
assert!(parse_increment_channel(&data).is_none());
}
#[test]
fn projection_boundary_rejects_invalid_unread_counts() {
for key in ["unreadCount", "unread_count"] {
for value in [
serde_json::json!(-1),
serde_json::json!("bad"),
serde_json::Value::Null,
] {
let mut data = serde_json::json!({"id": "bi643wys7fgy9pbai9wjfusnfa"});
data[key] = value;
assert!(parse_increment_channel(&data).is_none(), "{data}");
}
}
}
}