Skip to main content

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}