Skip to main content

servers/
batcher.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Logical-table and ordinary-table batching implementations.
16
17mod flow_notifier;
18mod flow_sender;
19pub mod logical_table;
20pub mod table;
21
22#[cfg(test)]
23mod test_util;
24
25use serde::{Deserialize, Serialize};
26
27/// Controls whether batching waits for storage before replying to the client.
28const PENDING_ROWS_BATCH_SYNC_ENV: &str = "PENDING_ROWS_BATCH_SYNC";
29
30/// Returns whether pending-row batch submissions wait for the flush result
31/// before replying to the client (synchronous mode), controlled by the
32/// `PENDING_ROWS_BATCH_SYNC` environment variable and defaulting to `true`.
33///
34/// Callers that reason about how long a remote write request may block (e.g.
35/// the frontend HTTP timeout fallback) must consult this instead of
36/// duplicating the env lookup.
37pub fn pending_rows_batch_sync_enabled() -> bool {
38    std::env::var(PENDING_ROWS_BATCH_SYNC_ENV)
39        .ok()
40        .as_deref()
41        .and_then(|v| v.parse::<bool>().ok())
42        .unwrap_or(true)
43}
44
45/// Ingestion protocols that can opt into pending-row batching.
46#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
47#[serde(rename_all = "snake_case")]
48pub enum BatchingProtocol {
49    Prom,
50    Influxdb,
51    Opentsdb,
52    Otlp,
53    Logs,
54    Loki,
55    Splunk,
56    Elasticsearch,
57    HttpSql,
58    Mysql,
59    Postgres,
60}
61
62#[cfg(test)]
63mod tests {
64    use super::*;
65
66    #[test]
67    fn test_protocol_names_reject_unknown_values() {
68        assert_eq!(
69            serde_json::from_str::<BatchingProtocol>("\"prom\"").unwrap(),
70            BatchingProtocol::Prom
71        );
72        for name in ["sql", "unknown"] {
73            assert!(serde_json::from_str::<BatchingProtocol>(&format!("\"{name}\"")).is_err());
74        }
75        assert_eq!(
76            serde_json::from_str::<BatchingProtocol>("\"http_sql\"").unwrap(),
77            BatchingProtocol::HttpSql
78        );
79    }
80}