surreal_sync_runtime/pipeline/external/wire.rs
1//! NDJSON wire headers for external transform requests/responses.
2
3use serde::{Deserialize, Serialize};
4
5/// Item payload kind for an External NDJSON exchange.
6///
7/// Workers that only handle row CDC may ignore `kind` (default / omitted =
8/// [`WireItemKind::Change`]). Relation-aware workers must honor `kind` so
9/// relation batches are not silently treated as row changes.
10#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
11#[serde(rename_all = "snake_case")]
12pub enum WireItemKind {
13 /// Incremental row [`surreal_sync_core::Change`] items.
14 #[default]
15 Change,
16 /// Full-sync / snapshot [`surreal_sync_core::Row`] items.
17 Row,
18 /// Incremental relation [`surreal_sync_core::RelationChange`] items.
19 RelationChange,
20 /// Full-sync [`surreal_sync_core::Relation`] items.
21 Relation,
22}
23
24/// Request header line: `{"batch_id","count"}` plus optional `kind`.
25#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
26pub struct RequestHeader {
27 pub batch_id: u64,
28 pub count: usize,
29 /// Payload kind. Omitted / default = [`WireItemKind::Change`] for back-compat.
30 #[serde(default, skip_serializing_if = "is_default_kind")]
31 pub kind: WireItemKind,
32}
33
34fn is_default_kind(kind: &WireItemKind) -> bool {
35 *kind == WireItemKind::Change
36}
37
38/// Response header line: must echo `batch_id`; either `count` + items or `error`.
39#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
40pub struct ResponseHeader {
41 pub batch_id: u64,
42 #[serde(default, skip_serializing_if = "Option::is_none")]
43 pub count: Option<usize>,
44 #[serde(default, skip_serializing_if = "Option::is_none")]
45 pub error: Option<String>,
46}