Skip to main content

servers/http/
header.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
15use std::collections::HashMap;
16use std::sync::Arc;
17
18use common_plugins::GREPTIME_EXEC_PREFIX;
19use datafusion::physical_plan::ExecutionPlan;
20use datafusion::physical_plan::metrics::MetricValue;
21use headers::{Header, HeaderName, HeaderValue};
22use hyper::HeaderMap;
23use serde_json::Value;
24
25pub mod constants {
26    // New HTTP headers would better distinguish use cases among:
27    // * GreptimeDB
28    // * GreptimeCloud
29    // * ...
30    //
31    // And thus trying to use:
32    // * x-greptime-db-xxx
33    // * x-greptime-cloud-xxx
34    //
35    // ... accordingly
36    //
37    // Most of the headers are for GreptimeDB and thus using `x-greptime-db-` as prefix.
38    // Only use `x-greptime-cloud` when it's intentionally used by GreptimeCloud.
39
40    // LEGACY HEADERS - KEEP IT UNMODIFIED
41    pub const GREPTIME_DB_HEADER_FORMAT: &str = "x-greptime-format";
42    pub const GREPTIME_DB_HEADER_TIMEOUT: &str = "x-greptime-timeout";
43    pub const GREPTIME_DB_HEADER_EXECUTION_TIME: &str = "x-greptime-execution-time";
44    pub const GREPTIME_DB_HEADER_METRICS: &str = "x-greptime-metrics";
45    pub const GREPTIME_DB_HEADER_NAME: &str = "x-greptime-db-name";
46    pub const GREPTIME_DB_HEADER_READ_PREFERENCE: &str = "x-greptime-read-preference";
47    pub const GREPTIME_INSERT_SKIP_WAL_HEADER_NAME: &str = "x-greptime-insert-skip-wal";
48    pub const GREPTIME_TIMEZONE_HEADER_NAME: &str = "x-greptime-timezone";
49    pub const GREPTIME_DB_HEADER_ERROR_CODE: &str = common_error::GREPTIME_DB_HEADER_ERROR_CODE;
50
51    // Deprecated: pipeline is also used with trace, so we remove log from it.
52    pub const GREPTIME_LOG_PIPELINE_NAME_HEADER_NAME: &str = "x-greptime-log-pipeline-name";
53    pub const GREPTIME_LOG_PIPELINE_VERSION_HEADER_NAME: &str = "x-greptime-log-pipeline-version";
54
55    // More generic pipeline header name
56    pub const GREPTIME_PIPELINE_NAME_HEADER_NAME: &str = "x-greptime-pipeline-name";
57    pub const GREPTIME_PIPELINE_VERSION_HEADER_NAME: &str = "x-greptime-pipeline-version";
58
59    pub const GREPTIME_LOG_TABLE_NAME_HEADER_NAME: &str = "x-greptime-log-table-name";
60    pub const GREPTIME_LOG_EXTRACT_KEYS_HEADER_NAME: &str = "x-greptime-log-extract-keys";
61    pub const GREPTIME_TRACE_TABLE_NAME_HEADER_NAME: &str = "x-greptime-trace-table-name";
62
63    // OTLP headers
64    pub const GREPTIME_OTLP_METRIC_PROMOTE_ALL_RESOURCE_ATTRS_HEADER_NAME: &str =
65        "x-greptime-otlp-metric-promote-all-resource-attrs";
66    pub const GREPTIME_OTLP_METRIC_PROMOTE_RESOURCE_ATTRS_HEADER_NAME: &str =
67        "x-greptime-otlp-metric-promote-resource-attrs";
68    pub const GREPTIME_OTLP_METRIC_IGNORE_RESOURCE_ATTRS_HEADER_NAME: &str =
69        "x-greptime-otlp-metric-ignore-resource-attrs";
70    pub const GREPTIME_OTLP_METRIC_PROMOTE_SCOPE_ATTRS_HEADER_NAME: &str =
71        "x-greptime-otlp-metric-promote-scope-attrs";
72    pub const GREPTIME_OTLP_METRIC_TRANSLATION_STRATEGY_HEADER_NAME: &str =
73        "x-greptime-otlp-metric-translation-strategy";
74
75    /// The header key that contains the pipeline params.
76    pub const GREPTIME_PIPELINE_PARAMS_HEADER: &str = "x-greptime-pipeline-params";
77}
78
79pub static GREPTIME_DB_HEADER_FORMAT: HeaderName =
80    HeaderName::from_static(constants::GREPTIME_DB_HEADER_FORMAT);
81pub static GREPTIME_DB_HEADER_EXECUTION_TIME: HeaderName =
82    HeaderName::from_static(constants::GREPTIME_DB_HEADER_EXECUTION_TIME);
83pub static GREPTIME_DB_HEADER_METRICS: HeaderName =
84    HeaderName::from_static(constants::GREPTIME_DB_HEADER_METRICS);
85
86/// Header key of `db-name`. Example format of the header value is `greptime-public`.
87pub static GREPTIME_DB_HEADER_NAME: HeaderName =
88    HeaderName::from_static(constants::GREPTIME_DB_HEADER_NAME);
89
90/// Header key of query specific timezone. Example format of the header value is `Asia/Shanghai` or `+08:00`.
91pub static GREPTIME_TIMEZONE_HEADER_NAME: HeaderName =
92    HeaderName::from_static(constants::GREPTIME_TIMEZONE_HEADER_NAME);
93
94/// Header key of query specific read preference. Example format of the header value is `leader`.
95pub static GREPTIME_DB_HEADER_READ_PREFERENCE: HeaderName =
96    HeaderName::from_static(constants::GREPTIME_DB_HEADER_READ_PREFERENCE);
97
98/// Request-level WAL policy, independent of the table-level skip_wal option.
99pub static GREPTIME_INSERT_SKIP_WAL_HEADER_NAME: HeaderName =
100    HeaderName::from_static(constants::GREPTIME_INSERT_SKIP_WAL_HEADER_NAME);
101
102pub static CONTENT_TYPE_PROTOBUF_STR: &str = "application/x-protobuf";
103pub static CONTENT_TYPE_PROTOBUF: HeaderValue = HeaderValue::from_static(CONTENT_TYPE_PROTOBUF_STR);
104pub static CONTENT_ENCODING_SNAPPY: HeaderValue = HeaderValue::from_static("snappy");
105
106pub static CONTENT_TYPE_NDJSON_STR: &str = "application/x-ndjson";
107pub static CONTENT_TYPE_NDJSON_SUBTYPE_STR: &str = "x-ndjson";
108
109pub struct GreptimeDbName(Option<String>);
110
111impl Header for GreptimeDbName {
112    fn name() -> &'static HeaderName {
113        &GREPTIME_DB_HEADER_NAME
114    }
115
116    fn decode<'i, I>(values: &mut I) -> Result<Self, headers::Error>
117    where
118        Self: Sized,
119        I: Iterator<Item = &'i HeaderValue>,
120    {
121        if let Some(value) = values.next() {
122            let str_value = value.to_str().map_err(|_| headers::Error::invalid())?;
123            Ok(Self(Some(str_value.to_owned())))
124        } else {
125            Ok(Self(None))
126        }
127    }
128
129    fn encode<E: Extend<HeaderValue>>(&self, values: &mut E) {
130        if let Some(name) = &self.0
131            && let Ok(value) = HeaderValue::from_str(name)
132        {
133            values.extend(std::iter::once(value));
134        }
135    }
136}
137
138impl GreptimeDbName {
139    pub fn value(&self) -> Option<&String> {
140        self.0.as_ref()
141    }
142}
143
144// collect write
145pub fn write_cost_header_map(cost: usize) -> HeaderMap {
146    let mut header_map = HeaderMap::new();
147    if cost > 0 {
148        let mut map: HashMap<String, Value> = HashMap::new();
149        map.insert(
150            common_plugins::GREPTIME_EXEC_WRITE_COST.to_string(),
151            Value::from(cost),
152        );
153        let _ = serde_json::to_string(&map)
154            .ok()
155            .and_then(|s| HeaderValue::from_str(&s).ok())
156            .and_then(|v| header_map.insert(&GREPTIME_DB_HEADER_METRICS, v));
157    }
158    header_map
159}
160
161fn collect_into_maps(name: &str, value: u64, maps: &mut [&mut HashMap<String, u64>]) {
162    if name.starts_with(GREPTIME_EXEC_PREFIX) && value > 0 {
163        maps.iter_mut().for_each(|map| {
164            map.entry(name.to_string())
165                .and_modify(|v| *v += value)
166                .or_insert(value);
167        });
168    }
169}
170
171pub fn collect_plan_metrics(plan: &Arc<dyn ExecutionPlan>, maps: &mut [&mut HashMap<String, u64>]) {
172    if let Some(m) = plan.metrics() {
173        m.iter().for_each(|m| match m.value() {
174            MetricValue::Count { name, count } => {
175                collect_into_maps(name, count.value() as u64, maps);
176            }
177            MetricValue::Gauge { name, gauge } => {
178                collect_into_maps(name, gauge.value() as u64, maps);
179            }
180            MetricValue::Time { name, time } if name.starts_with(GREPTIME_EXEC_PREFIX) => {
181                // override
182                maps.iter_mut().for_each(|map| {
183                    map.insert(name.to_string(), time.value() as u64);
184                });
185            }
186            _ => {}
187        });
188    }
189
190    for c in plan.children() {
191        collect_plan_metrics(c, maps);
192    }
193}