Skip to main content

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}