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),
};
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),
};
let kind = entry["kind"]
.as_str()
.ok_or_else(|| ImError::Parse("sync entry missing 'kind'".to_string()))?;
match kind {
"no_change" => Ok(SyncResponse::NoChange),
"events" => {
let mut events = parse_channel_events(entry, channel_id)?;
let messages = parse_messages_map(entry);
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,
})
}
"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_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())
.or_else(|| ev["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);
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) -> std::collections::HashMap<String, PostFields> {
let mut map = std::collections::HashMap::new();
let Some(obj) = entry.get("messages").and_then(|m| m.as_object()) else {
return map;
};
map.reserve(obj.len());
for (msg_id, post) in obj {
map.insert(msg_id.clone(), extract_post_fields(post));
}
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,
});
}