nodedb_types/sync/wire/timeseries.rs
1// SPDX-License-Identifier: Apache-2.0
2
3//! Timeseries ingest + definition-sync messages.
4
5use std::collections::HashMap;
6
7use serde::{Deserialize, Serialize};
8
9use crate::sync::wire::ack_status::AckStatus;
10
11/// Timeseries metric batch push (client → server, 0x40).
12#[derive(
13 Debug, Clone, Serialize, Deserialize, zerompk::ToMessagePack, zerompk::FromMessagePack,
14)]
15pub struct TimeseriesPushMsg {
16 /// Source Lite instance ID (UUID v7).
17 pub lite_id: String,
18 /// Collection name.
19 pub collection: String,
20 /// Gorilla-encoded timestamp block.
21 pub ts_block: Vec<u8>,
22 /// Gorilla-encoded value block.
23 pub val_block: Vec<u8>,
24 /// Raw LE u64 series ID block.
25 pub series_block: Vec<u8>,
26 /// Number of samples in this batch.
27 pub sample_count: u64,
28 /// Min timestamp in this batch.
29 pub min_ts: i64,
30 /// Max timestamp in this batch.
31 pub max_ts: i64,
32 /// Per-series sync watermark: highest LSN already synced for each series.
33 /// Only samples after these watermarks are included.
34 pub watermarks: HashMap<u64, u64>,
35 /// Stable identity of the originating producer. 0 for legacy clients.
36 #[serde(default)]
37 pub producer_id: u64,
38 /// Monotonic epoch counter incremented on every producer restart.
39 #[serde(default)]
40 pub epoch: u64,
41 /// Per-stream monotonic sequence number within the epoch.
42 #[serde(default)]
43 pub seq: u64,
44}
45
46/// Timeseries push acknowledgment (server → client, 0x41).
47#[derive(
48 Debug, Clone, Serialize, Deserialize, zerompk::ToMessagePack, zerompk::FromMessagePack,
49)]
50pub struct TimeseriesAckMsg {
51 /// Collection acknowledged.
52 pub collection: String,
53 /// Number of samples accepted.
54 pub accepted: u64,
55 /// Number of samples rejected (duplicates, out-of-retention, etc.)
56 pub rejected: u64,
57 /// Server-assigned LSN for this batch (used as sync watermark).
58 pub lsn: u64,
59 /// Highest sequence number from this producer that has been durably applied.
60 #[serde(default)]
61 pub applied_seq: u64,
62 /// Idempotency outcome of the acknowledged message.
63 #[serde(default)]
64 pub status: AckStatus,
65}
66
67/// Definition sync message (server → client, 0x70).
68///
69/// Carries function/trigger/procedure definitions from Origin to Lite.
70/// Sent when definitions are created, modified, or dropped on Origin.
71#[derive(
72 Debug, Clone, Serialize, Deserialize, zerompk::ToMessagePack, zerompk::FromMessagePack,
73)]
74pub struct DefinitionSyncMsg {
75 /// Type of definition: "function", "trigger", "procedure".
76 pub definition_type: String,
77 /// The definition name.
78 pub name: String,
79 /// Action: "put" (create/replace) or "delete" (drop).
80 pub action: String,
81 /// Serialized definition body (JSON). Empty for "delete" actions.
82 pub payload: Vec<u8>,
83}