1use std::collections::BTreeMap;
16
17use api::v1::ColumnDataType;
18use api::v1::value::ValueData;
19use ordered_float::NotNan;
20use snafu::{OptionExt, ResultExt};
21use vrl::prelude::Bytes;
22use vrl::value::{KeyString, Value as VrlValue};
23
24use crate::error::{
25 FloatIsNanSnafu, Result, ValueDefaultValueUnsupportedSnafu, ValueInvalidResolutionSnafu,
26 ValueParseBooleanSnafu, ValueParseFloatSnafu, ValueParseIntSnafu, ValueParseTypeSnafu,
27 ValueUnsupportedYamlTypeSnafu, ValueYamlKeyMustBeStringSnafu,
28};
29
30pub(crate) const NANOSECOND_RESOLUTION: &str = "nanosecond";
31pub(crate) const NANO_RESOLUTION: &str = "nano";
32pub(crate) const NS_RESOLUTION: &str = "ns";
33pub(crate) const MICROSECOND_RESOLUTION: &str = "microsecond";
34pub(crate) const MICRO_RESOLUTION: &str = "micro";
35pub(crate) const US_RESOLUTION: &str = "us";
36pub(crate) const MILLISECOND_RESOLUTION: &str = "millisecond";
37pub(crate) const MILLI_RESOLUTION: &str = "milli";
38pub(crate) const MS_RESOLUTION: &str = "ms";
39pub(crate) const SECOND_RESOLUTION: &str = "second";
40pub(crate) const SEC_RESOLUTION: &str = "sec";
41pub(crate) const S_RESOLUTION: &str = "s";
42
43pub(crate) const VALID_RESOLUTIONS: [&str; 12] = [
44 NANOSECOND_RESOLUTION,
45 NANO_RESOLUTION,
46 NS_RESOLUTION,
47 MICROSECOND_RESOLUTION,
48 MICRO_RESOLUTION,
49 US_RESOLUTION,
50 MILLISECOND_RESOLUTION,
51 MILLI_RESOLUTION,
52 MS_RESOLUTION,
53 SECOND_RESOLUTION,
54 SEC_RESOLUTION,
55 S_RESOLUTION,
56];
57
58pub fn parse_str_type(t: &str) -> Result<ColumnDataType> {
59 let mut parts = t.splitn(2, ',');
60 let head = parts.next().unwrap_or_default();
61 let tail = parts.next().map(|s| s.trim().to_string());
62 match head.to_lowercase().as_str() {
63 "int8" => Ok(ColumnDataType::Int8),
64 "int16" => Ok(ColumnDataType::Int16),
65 "int32" => Ok(ColumnDataType::Int32),
66 "int64" => Ok(ColumnDataType::Int64),
67
68 "uint8" => Ok(ColumnDataType::Uint8),
69 "uint16" => Ok(ColumnDataType::Uint16),
70 "uint32" => Ok(ColumnDataType::Uint32),
71 "uint64" => Ok(ColumnDataType::Uint64),
72
73 "float32" => Ok(ColumnDataType::Float32),
74 "float64" => Ok(ColumnDataType::Float64),
75
76 "boolean" => Ok(ColumnDataType::Boolean),
77 "string" => Ok(ColumnDataType::String),
78
79 "timestamp" | "epoch" | "time" => match tail {
80 Some(resolution) if !resolution.is_empty() => match resolution.as_str() {
81 NANOSECOND_RESOLUTION | NANO_RESOLUTION | NS_RESOLUTION => {
82 Ok(ColumnDataType::TimestampNanosecond)
83 }
84 MICROSECOND_RESOLUTION | MICRO_RESOLUTION | US_RESOLUTION => {
85 Ok(ColumnDataType::TimestampMicrosecond)
86 }
87 MILLISECOND_RESOLUTION | MILLI_RESOLUTION | MS_RESOLUTION => {
88 Ok(ColumnDataType::TimestampMillisecond)
89 }
90 SECOND_RESOLUTION | SEC_RESOLUTION | S_RESOLUTION => {
91 Ok(ColumnDataType::TimestampSecond)
92 }
93 _ => ValueInvalidResolutionSnafu {
94 resolution,
95 valid_resolution: VALID_RESOLUTIONS.join(","),
96 }
97 .fail(),
98 },
99 _ => Ok(ColumnDataType::TimestampNanosecond),
100 },
101
102 "json" => Ok(ColumnDataType::Binary),
105 "json2" => Ok(ColumnDataType::Json),
106
107 _ => ValueParseTypeSnafu { t }.fail(),
108 }
109}
110
111pub fn parse_str_value(type_: &ColumnDataType, v: &str) -> Result<ValueData> {
112 match type_ {
113 ColumnDataType::Int8 => v
114 .parse::<i8>()
115 .map(|v| ValueData::I8Value(v as i32))
116 .context(ValueParseIntSnafu { ty: "int8", v }),
117 ColumnDataType::Int16 => v
118 .parse::<i16>()
119 .map(|v| ValueData::I16Value(v as i32))
120 .context(ValueParseIntSnafu { ty: "int16", v }),
121 ColumnDataType::Int32 => v
122 .parse::<i32>()
123 .map(ValueData::I32Value)
124 .context(ValueParseIntSnafu { ty: "int32", v }),
125 ColumnDataType::Int64 => v
126 .parse::<i64>()
127 .map(ValueData::I64Value)
128 .context(ValueParseIntSnafu { ty: "int64", v }),
129
130 ColumnDataType::Uint8 => v
131 .parse::<u8>()
132 .map(|v| ValueData::U8Value(v as u32))
133 .context(ValueParseIntSnafu { ty: "uint8", v }),
134 ColumnDataType::Uint16 => v
135 .parse::<u16>()
136 .map(|v| ValueData::U16Value(v as u32))
137 .context(ValueParseIntSnafu { ty: "uint16", v }),
138 ColumnDataType::Uint32 => v
139 .parse::<u32>()
140 .map(ValueData::U32Value)
141 .context(ValueParseIntSnafu { ty: "uint32", v }),
142 ColumnDataType::Uint64 => v
143 .parse::<u64>()
144 .map(ValueData::U64Value)
145 .context(ValueParseIntSnafu { ty: "uint64", v }),
146
147 ColumnDataType::Float32 => v
148 .parse::<f32>()
149 .map(ValueData::F32Value)
150 .context(ValueParseFloatSnafu { ty: "float32", v }),
151 ColumnDataType::Float64 => v
152 .parse::<f64>()
153 .map(ValueData::F64Value)
154 .context(ValueParseFloatSnafu { ty: "float64", v }),
155
156 ColumnDataType::Boolean => v
157 .parse::<bool>()
158 .map(ValueData::BoolValue)
159 .context(ValueParseBooleanSnafu { ty: "boolean", v }),
160 ColumnDataType::String => Ok(ValueData::StringValue(v.to_string())),
161
162 _ => ValueDefaultValueUnsupportedSnafu {
163 value: format!("{:?}", type_),
164 }
165 .fail(),
166 }
167}
168
169pub fn yaml_to_vrl_value(v: &yaml_rust::Yaml) -> Result<VrlValue> {
170 match v {
171 yaml_rust::Yaml::Null => Ok(VrlValue::Null),
172 yaml_rust::Yaml::Boolean(v) => Ok(VrlValue::Boolean(*v)),
173 yaml_rust::Yaml::Integer(v) => Ok(VrlValue::Integer(*v)),
174 yaml_rust::Yaml::Real(v) => {
175 let f = v
176 .parse::<f64>()
177 .context(ValueParseFloatSnafu { ty: "float64", v })?;
178 NotNan::new(f).map(VrlValue::Float).context(FloatIsNanSnafu)
179 }
180 yaml_rust::Yaml::String(v) => Ok(VrlValue::Bytes(Bytes::from(v.clone()))),
181 yaml_rust::Yaml::Array(arr) => {
182 let mut values = vec![];
183 for v in arr {
184 values.push(yaml_to_vrl_value(v)?);
185 }
186 Ok(VrlValue::Array(values))
187 }
188 yaml_rust::Yaml::Hash(v) => {
189 let mut values = BTreeMap::new();
190 for (k, v) in v {
191 let key = k
192 .as_str()
193 .with_context(|| ValueYamlKeyMustBeStringSnafu { value: v.clone() })?;
194 values.insert(KeyString::from(key), yaml_to_vrl_value(v)?);
195 }
196 Ok(VrlValue::Object(values))
197 }
198 _ => ValueUnsupportedYamlTypeSnafu { value: v.clone() }.fail(),
199 }
200}