use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
pub type Version = i64;
pub const MAX_PUSH_CHANGES: usize = 1000;
pub const MAX_PULL_LIMIT: i64 = 1000;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum Op {
Upsert,
Delete,
}
mod payload_presence {
use serde::{Deserialize, Deserializer, Serialize, Serializer};
use serde_json::Value;
#[allow(clippy::ref_option)]
pub fn serialize<S: Serializer>(
payload: &Option<Value>,
serializer: S,
) -> Result<S::Ok, S::Error> {
payload.serialize(serializer)
}
pub fn deserialize<'de, D: Deserializer<'de>>(
deserializer: D,
) -> Result<Option<Value>, D::Error> {
Value::deserialize(deserializer).map(Some)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Change {
pub change_id: String,
pub collection: String,
pub pk: String,
pub op: Op,
#[serde(
default,
skip_serializing_if = "Option::is_none",
with = "payload_presence"
)]
pub payload: Option<serde_json::Value>,
pub base_version: Version,
pub updated_at: DateTime<Utc>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PushRequest {
pub device_id: String,
pub changes: Vec<Change>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RemoteRow {
pub collection: String,
pub pk: String,
#[serde(
default,
skip_serializing_if = "Option::is_none",
with = "payload_presence"
)]
pub payload: Option<serde_json::Value>,
pub version: Version,
pub deleted: bool,
pub updated_at: DateTime<Utc>,
pub device_id: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "status", rename_all = "snake_case")]
pub enum ChangeOutcome {
Applied {
version: Version,
},
AlreadyApplied {
version: Version,
},
Resolved {
row: RemoteRow,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PushResponse {
pub outcomes: Vec<ChangeOutcome>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub struct PullQuery {
#[serde(default)]
pub cursor: Version,
#[serde(default = "default_pull_limit")]
pub limit: i64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub session: Option<Version>,
}
const fn default_pull_limit() -> i64 {
500
}
impl PullQuery {
#[must_use]
pub fn session_start(&self) -> Version {
self.session.unwrap_or(self.cursor)
}
}
impl Default for PullQuery {
fn default() -> Self {
Self {
cursor: 0,
limit: default_pull_limit(),
session: None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "status", rename_all = "snake_case")]
pub enum PullResponse {
Ok {
rows: Vec<RemoteRow>,
next_cursor: Version,
tombstone_horizon: Version,
},
FullResyncRequired {
tombstone_horizon: Version,
},
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::TimeZone;
fn sample_change() -> Change {
Change {
change_id: "11111111-1111-4111-8111-111111111111".into(),
collection: "notes".into(),
pk: "n1".into(),
op: Op::Upsert,
payload: Some(serde_json::json!({"title": "hi"})),
base_version: 3,
updated_at: Utc.with_ymd_and_hms(2026, 7, 7, 12, 0, 0).unwrap(),
}
}
#[test]
fn push_request_round_trips() {
let request = PushRequest {
device_id: "device-a".into(),
changes: vec![sample_change()],
};
let json = serde_json::to_string(&request).unwrap();
let back: PushRequest = serde_json::from_str(&json).unwrap();
assert_eq!(back, request);
assert!(json.contains(r#""op":"upsert""#), "snake_case ops: {json}");
}
#[test]
fn delete_change_omits_payload() {
let mut change = sample_change();
change.op = Op::Delete;
change.payload = None;
let json = serde_json::to_string(&change).unwrap();
assert!(
!json.contains("payload"),
"no payload key for deletes: {json}"
);
let back: Change = serde_json::from_str(&json).unwrap();
assert_eq!(back, change);
}
#[test]
fn payload_presence_distinguishes_absent_null_and_value() {
let mut change = sample_change();
change.op = Op::Delete;
change.payload = None;
let json = serde_json::to_string(&change).unwrap();
assert!(
!json.contains("payload"),
"absent must omit the key: {json}"
);
assert_eq!(serde_json::from_str::<Change>(&json).unwrap(), change);
let mut change = sample_change();
change.payload = Some(serde_json::Value::Null);
let json = serde_json::to_string(&change).unwrap();
assert!(
json.contains(r#""payload":null"#),
"present null must stay on the wire: {json}"
);
assert_eq!(serde_json::from_str::<Change>(&json).unwrap(), change);
let change = sample_change();
let json = serde_json::to_string(&change).unwrap();
assert_eq!(serde_json::from_str::<Change>(&json).unwrap(), change);
}
#[test]
fn remote_row_payload_presence_round_trips() {
let row = |payload: Option<serde_json::Value>, deleted: bool| RemoteRow {
collection: "notes".to_owned(),
pk: "n1".to_owned(),
payload,
version: 7,
deleted,
updated_at: chrono::Utc::now(),
device_id: "d".to_owned(),
};
let tombstone = row(None, true);
let json = serde_json::to_string(&tombstone).unwrap();
assert!(!json.contains("payload"), "{json}");
assert_eq!(serde_json::from_str::<RemoteRow>(&json).unwrap(), tombstone);
let null_row = row(Some(serde_json::Value::Null), false);
let json = serde_json::to_string(&null_row).unwrap();
assert!(json.contains(r#""payload":null"#), "{json}");
assert_eq!(serde_json::from_str::<RemoteRow>(&json).unwrap(), null_row);
let snapshot = serde_json::to_value(&null_row).unwrap();
assert_eq!(
serde_json::from_value::<RemoteRow>(snapshot).unwrap(),
null_row
);
let live = row(Some(serde_json::json!({"a": 1})), false);
let json = serde_json::to_string(&live).unwrap();
assert_eq!(serde_json::from_str::<RemoteRow>(&json).unwrap(), live);
}
#[test]
fn change_outcomes_are_status_tagged() {
let applied = serde_json::to_value(ChangeOutcome::Applied { version: 7 }).unwrap();
assert_eq!(applied["status"], "applied");
assert_eq!(applied["version"], 7);
let deduped = serde_json::to_value(ChangeOutcome::AlreadyApplied { version: 7 }).unwrap();
assert_eq!(deduped["status"], "already_applied");
assert_eq!(
deduped["version"], 7,
"already_applied must carry the originally assigned version"
);
}
#[test]
fn pull_response_round_trips_both_variants() {
let ok = PullResponse::Ok {
rows: vec![RemoteRow {
collection: "notes".into(),
pk: "n1".into(),
payload: None,
version: 9,
deleted: true,
updated_at: Utc.with_ymd_and_hms(2026, 7, 7, 12, 0, 0).unwrap(),
device_id: "device-a".into(),
}],
next_cursor: 9,
tombstone_horizon: 4,
};
let json = serde_json::to_string(&ok).unwrap();
assert!(json.contains(r#""status":"ok""#));
assert_eq!(serde_json::from_str::<PullResponse>(&json).unwrap(), ok);
let resync = PullResponse::FullResyncRequired {
tombstone_horizon: 4,
};
let json = serde_json::to_string(&resync).unwrap();
assert!(json.contains(r#""status":"full_resync_required""#));
assert_eq!(serde_json::from_str::<PullResponse>(&json).unwrap(), resync);
}
#[test]
fn pull_query_defaults() {
let query: PullQuery = serde_urlencoded::from_str("").unwrap();
assert_eq!(query, PullQuery::default());
assert_eq!(query.cursor, 0);
assert_eq!(query.limit, 500);
assert_eq!(query.session, None);
assert_eq!(
query.session_start(),
0,
"no session marker means the page cursor is the session start"
);
let query: PullQuery = serde_urlencoded::from_str("cursor=12&limit=50").unwrap();
assert_eq!(query.cursor, 12);
assert_eq!(query.limit, 50);
assert_eq!(query.session_start(), 12);
let query: PullQuery = serde_urlencoded::from_str("cursor=12&limit=50&session=3").unwrap();
assert_eq!(query.session, Some(3));
assert_eq!(
query.session_start(),
3,
"an explicit session-start cursor wins over the page cursor"
);
}
}