use crate::error::ImError;
use crate::state::{ChannelId, Seq};
use crate::sync_session::{ChannelSnapshot, EventEnvelope, EventKind, PostFields, SyncResponse};
use super::fields::extract_post_fields;
use super::kind::parse_event_kind;
pub const SYNC_EVENTS_PER_ENTRY: usize = 500;
pub fn parse_sync_response(bytes: &[u8], channel_id: ChannelId) -> Result<SyncResponse, ImError> {
let v: serde_json::Value = serde_json::from_slice(bytes)
.map_err(|e| ImError::Parse(format!("sync response invalid JSON: {}", e)))?;
if let Some(status) = v.get("status").and_then(|s| s.as_str()) {
if status != "SUCCESS" {
return Err(ImError::Parse(format!(
"sync response status not SUCCESS: {}",
status
)));
}
}
let entries = v
.get("data")
.and_then(|d| d.get("entries"))
.or_else(|| v.get("entries"))
.and_then(|e| e.as_array());
let entries = match entries {
Some(e) => e,
None => {
return Ok(SyncResponse::NoChange {
next_seq: Seq(0),
persona: Default::default(),
})
}
};
let entry = entries.iter().find(|e| {
let cid = e
.get("channelId")
.or_else(|| e.get("channel_id"))
.and_then(|c| c.as_str());
cid == Some(channel_id.as_str())
});
let entry = match entry {
Some(e) => e,
None => {
return Ok(SyncResponse::NoChange {
next_seq: Seq(0),
persona: Default::default(),
})
}
};
let kind = entry["kind"]
.as_str()
.ok_or_else(|| ImError::Parse("sync entry missing 'kind'".to_string()))?;
match kind {
"no_change" => Ok(SyncResponse::NoChange {
next_seq: Seq(entry["nextSeq"]
.as_u64()
.or_else(|| entry["next_seq"].as_u64())
.unwrap_or(0)),
persona: parse_sync_persona(entry)?,
}),
"events" => {
let canonical = entry
.get("streamEvents")
.or_else(|| entry.get("stream_events"));
let mut events = if canonical.is_some() {
parse_stream_events(entry, channel_id)?
} else {
parse_channel_events(entry, channel_id)?
};
let messages = parse_messages_map(entry)?;
if canonical.is_some()
&& events.iter().any(|event| {
!event.redacted
&& event
.msg_id
.as_deref()
.is_some_and(|msg_id| !messages.contains_key(msg_id))
})
{
return Err(ImError::Parse(
"canonical stream event body missing from messages".to_string(),
));
}
let needs_continuation = events.len() >= SYNC_EVENTS_PER_ENTRY;
let max_event_seq = events.iter().map(|e| e.seq.0).max().unwrap_or(0);
let next_seq = Seq(entry["nextSeq"]
.as_u64()
.or_else(|| entry["next_seq"].as_u64())
.unwrap_or(max_event_seq));
sort_events_by_bucket(&mut events);
Ok(SyncResponse::Events {
events,
messages,
next_seq,
needs_continuation,
persona: parse_sync_persona(entry)?,
})
}
"too_long" => {
let reset_to = entry["resetTo"]
.as_u64()
.or_else(|| entry["reset_to"].as_u64())
.ok_or_else(|| ImError::Parse("too_long missing resetTo".to_string()))?;
Ok(SyncResponse::TooLong {
reset_to: Seq(reset_to),
})
}
"snapshot" => {
let reset_to = entry["resetTo"]
.as_u64()
.or_else(|| entry["reset_to"].as_u64())
.unwrap_or(0);
let messages = parse_channel_events(entry, channel_id)?;
if messages
.iter()
.any(|event| event.kind == EventKind::ChannelTerminalClosed)
{
return Err(ImError::Parse(
"terminal event is only valid in sync events entries".to_string(),
));
}
Ok(SyncResponse::Snapshot(ChannelSnapshot {
channel_id,
reset_to: Seq(reset_to),
messages,
}))
}
other => Err(ImError::Parse(format!(
"unknown sync entry kind: {}",
other
))),
}
}
fn parse_sync_persona(
entry: &serde_json::Value,
) -> Result<crate::sync_session::SyncPersona, ImError> {
let seq = |camel: &str, snake: &str| {
entry
.get(camel)
.or_else(|| entry.get(snake))
.and_then(serde_json::Value::as_u64)
.map(Seq)
};
let persona = crate::sync_session::SyncPersona {
membership_state: entry
.get("membershipState")
.or_else(|| entry.get("membership_state"))
.and_then(serde_json::Value::as_str)
.unwrap_or_default()
.to_string(),
epoch_start_seq: seq("epochStartSeq", "epoch_start_seq"),
epoch_end_seq: seq("epochEndSeq", "epoch_end_seq"),
member_projection: entry
.get("memberProjection")
.or_else(|| entry.get("member_projection"))
.filter(|value| !value.is_null())
.cloned(),
};
match persona.membership_state.as_str() {
"" => {}
"active"
if persona.epoch_start_seq.is_some()
&& persona.epoch_end_seq.is_none()
&& persona.member_projection.is_some() => {}
"ended"
if persona.epoch_start_seq.is_some()
&& persona.epoch_end_seq.is_some()
&& persona.member_projection.is_some() => {}
"never" | "invalid"
if persona.epoch_start_seq.is_none()
&& persona.epoch_end_seq.is_none()
&& persona.member_projection.is_none() => {}
_ => {
return Err(ImError::Parse(
"invalid sync membership persona".to_string(),
))
}
}
Ok(persona)
}
fn parse_stream_events(
entry: &serde_json::Value,
fallback_channel: ChannelId,
) -> Result<Vec<EventEnvelope>, ImError> {
let raw_events = entry
.get("streamEvents")
.or_else(|| entry.get("stream_events"))
.and_then(serde_json::Value::as_array)
.ok_or_else(|| ImError::Parse("streamEvents must be an array".to_string()))?;
let mut events = Vec::with_capacity(raw_events.len());
for raw in raw_events {
let Some(event) = crate::ws::handlers::channel_stream_event::parse_stream_event(raw, "")?
else {
return Err(ImError::Parse("invalid canonical stream event".to_string()));
};
if event.channel_id != fallback_channel {
return Err(ImError::Parse(
"canonical stream event channel mismatch".to_string(),
));
}
events.push(event);
}
Ok(events)
}
fn parse_channel_events(
entry: &serde_json::Value,
fallback_channel: ChannelId,
) -> Result<Vec<EventEnvelope>, ImError> {
let raw_events = match entry.get("events").and_then(|e| e.as_array()) {
Some(e) => e,
None => return Ok(Vec::new()),
};
let mut events = Vec::with_capacity(raw_events.len());
for ev in raw_events {
if ev.get("eventType").and_then(serde_json::Value::as_u64) == Some(7) {
events.push(parse_terminal_event(ev, fallback_channel)?);
continue;
}
let channel_id = ev
.get("channelId")
.and_then(|c| c.as_str())
.and_then(ChannelId::from_str)
.unwrap_or(fallback_channel);
let seq = ev["eventSeq"]
.as_u64()
.or_else(|| ev["event_seq"].as_u64())
.ok_or_else(|| ImError::Parse("channel event missing eventSeq".to_string()))?;
let event_type = ev["eventType"]
.as_u64()
.or_else(|| ev["event_type"].as_u64())
.unwrap_or(1) as u8;
let kind = parse_event_kind(event_type)?;
let fields = extract_post_fields(ev);
crate::category_chain::post::validate_fields(&fields)?;
let msg_id = ev
.get("msgId")
.or_else(|| ev.get("msg_id"))
.and_then(|v| v.as_str())
.map(str::to_string);
let event_id = ev
.get("id")
.or_else(|| ev.get("eventId"))
.and_then(|v| v.as_str())
.map(str::to_string);
let actor_id = ev
.get("actorId")
.or_else(|| ev.get("actor_id"))
.and_then(|v| v.as_str())
.map(str::to_string);
let occurred_at = ev
.get("occurredAt")
.or_else(|| ev.get("createAt"))
.or_else(|| ev.get("created_at"))
.and_then(|v| v.as_i64())
.unwrap_or(0);
let event_payload = ev
.get("payload")
.filter(|payload| !payload.is_null())
.and_then(|payload| serde_json::to_string(payload).ok())
.unwrap_or_default();
events.push(
EventEnvelope::new(channel_id, Seq(seq), kind, fields)
.with_msg_id(msg_id)
.with_event_identity(event_id, actor_id, occurred_at, event_payload),
);
}
Ok(events)
}
fn parse_terminal_event(
event: &serde_json::Value,
expected_channel: ChannelId,
) -> Result<EventEnvelope, ImError> {
let object = event
.as_object()
.ok_or_else(|| ImError::Parse("terminal event must be an object".to_string()))?;
const REQUIRED: [&str; 5] = ["id", "channelId", "eventSeq", "eventType", "payload"];
if object.len() != REQUIRED.len() || REQUIRED.iter().any(|key| !object.contains_key(*key)) {
return Err(ImError::Parse(
"terminal event must contain only id,channelId,eventSeq,eventType,payload".to_string(),
));
}
let id = object
.get("id")
.and_then(serde_json::Value::as_str)
.filter(|id| !id.is_empty())
.ok_or_else(|| ImError::Parse("terminal event missing/invalid id".to_string()))?;
let channel_id = object
.get("channelId")
.and_then(serde_json::Value::as_str)
.and_then(ChannelId::from_str)
.ok_or_else(|| ImError::Parse("terminal event missing/invalid channelId".to_string()))?;
if channel_id != expected_channel {
return Err(ImError::Parse(
"terminal event channelId does not match sync entry".to_string(),
));
}
let seq = object
.get("eventSeq")
.and_then(serde_json::Value::as_u64)
.filter(|seq| *seq > 0 && *seq <= i64::MAX as u64)
.ok_or_else(|| ImError::Parse("terminal event missing/invalid eventSeq".to_string()))?;
if object.get("eventType").and_then(serde_json::Value::as_u64) != Some(7) {
return Err(ImError::Parse(
"terminal event eventType must be 7".to_string(),
));
}
let payload = object
.get("payload")
.and_then(serde_json::Value::as_object)
.ok_or_else(|| ImError::Parse("terminal event payload must be an object".to_string()))?;
if payload.len() != 1
|| payload.get("state").and_then(serde_json::Value::as_str) != Some("closed")
{
return Err(ImError::Parse(
"terminal event payload must be exactly {state:closed}".to_string(),
));
}
Ok(EventEnvelope::new(
channel_id,
Seq(seq),
EventKind::ChannelTerminalClosed,
PostFields::default(),
)
.with_event_identity(Some(id.to_string()), None, 0, String::new()))
}
fn parse_messages_map(
entry: &serde_json::Value,
) -> Result<std::collections::HashMap<String, PostFields>, ImError> {
let mut map = std::collections::HashMap::new();
let Some(obj) = entry.get("messages").and_then(|m| m.as_object()) else {
return Ok(map);
};
map.reserve(obj.len());
for (msg_id, post) in obj {
let fields = extract_post_fields(post);
crate::category_chain::post::validate_fields(&fields)?;
map.insert(msg_id.clone(), fields);
}
Ok(map)
}
pub fn sort_events_by_bucket(events: &mut Vec<EventEnvelope>) {
events.sort_by_key(|ev| match ev.kind {
EventKind::PostUpsert => 0u8,
EventKind::PostEdit => 1,
EventKind::PostRevoke => 2,
EventKind::PostRead => 3,
EventKind::ChannelTerminalClosed => 4,
EventKind::Other(_) => 5,
});
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn sync_event_rejects_delivery_seq_as_stream_seq() {
let channel_id = crate::state::test_channel_id(42);
let body = serde_json::json!({
"data": {"entries": [{
"channelId": channel_id.as_str(),
"kind": "events",
"events": [{"channelId": channel_id.as_str(), "seq": 9001, "eventType": 1}]
}]}
});
assert!(parse_sync_response(body.to_string().as_bytes(), channel_id).is_err());
}
#[test]
fn canonical_sync_uses_stream_events_and_absolute_projection() {
let channel_id = crate::state::test_channel_id(43);
let body = serde_json::json!({
"data": {"entries": [{
"channelId": channel_id.as_str(),
"kind": "events",
"events": [{"eventSeq": 77, "eventType": 1}],
"streamEvents": [{
"channelId": channel_id.as_str(),
"streamSeq": 1,
"eventType": 1,
"eventId": "event-1",
"effectId": "effect-1",
"msgId": "message-1",
"redacted": false
}],
"messages": {"message-1": {"id": "message-1", "channelId": channel_id.as_str(), "message": "hello"}},
"memberProjection": {"channelId": channel_id.as_str(), "userId": "viewer", "projectionRevision": 2, "effectId": "effect-1", "unreadCount": 1},
"membershipState": "active",
"epochStartSeq": 1,
"nextSeq": 1
}]}
});
let SyncResponse::Events {
events, persona, ..
} = parse_sync_response(body.to_string().as_bytes(), channel_id).unwrap()
else {
panic!("expected events");
};
assert_eq!(events.len(), 1);
assert_eq!(events[0].seq.0, 1);
assert_eq!(events[0].effect_id, "effect-1");
assert_eq!(persona.membership_state, "active");
assert_eq!(persona.epoch_start_seq, Some(Seq(1)));
assert!(persona.member_projection.is_some());
}
#[test]
fn canonical_full_event_without_body_fails_closed_but_redacted_marker_does_not() {
let channel_id = crate::state::test_channel_id(44);
let entry = |redacted| {
serde_json::json!({
"data": {"entries": [{
"channelId": channel_id.as_str(),
"kind": "events",
"streamEvents": [{
"channelId": channel_id.as_str(),
"streamSeq": 1,
"eventType": 1,
"msgId": "message-1",
"redacted": redacted
}],
"nextSeq": 1
}]}
})
};
assert!(parse_sync_response(entry(false).to_string().as_bytes(), channel_id).is_err());
let SyncResponse::Events { events, .. } =
parse_sync_response(entry(true).to_string().as_bytes(), channel_id).unwrap()
else {
panic!("expected redacted events");
};
assert!(events[0].redacted);
assert!(events[0].msg_id.is_none());
}
}