1use std::io::Read;
2use std::io::Write;
3
4use chrono::{DateTime, Utc};
5use nu_protocol::shell_error::generic::GenericError;
6use nu_protocol::{PipelineData, Record, ShellError, Span, Value};
7
8use crate::store::Frame;
9use crate::store::Store;
10
11#[allow(clippy::result_large_err)] pub fn topic_value_to_string(value: Value) -> Result<String, ShellError> {
16 let span = value.span();
17 match value {
18 Value::String { val, .. } => Ok(val),
19 Value::List { vals, .. } => {
20 let mut parts = Vec::with_capacity(vals.len());
21 for v in vals {
22 let span = v.span();
23 match v {
24 Value::String { val, .. } => parts.push(val),
25 other => {
26 return Err(ShellError::Generic(GenericError::new(
27 "Invalid topic",
28 format!(
29 "topic list elements must be strings, got {}",
30 other.get_type()
31 ),
32 span,
33 )))
34 }
35 }
36 }
37 Ok(parts.join(","))
38 }
39 other => Err(ShellError::Generic(GenericError::new(
40 "Invalid topic",
41 format!(
42 "topic must be a string or a list of strings, got {}",
43 other.get_type()
44 ),
45 span,
46 ))),
47 }
48}
49
50pub fn json_to_value(json: &serde_json::Value, span: Span) -> Value {
51 match json {
52 serde_json::Value::Null => Value::nothing(span),
53 serde_json::Value::Bool(b) => Value::bool(*b, span),
54 serde_json::Value::Number(n) => {
55 if let Some(i) = n.as_i64() {
56 Value::int(i, span)
57 } else if let Some(f) = n.as_f64() {
58 Value::float(f, span)
59 } else {
60 Value::string(n.to_string(), span)
61 }
62 }
63 serde_json::Value::String(s) => Value::string(s, span),
64 serde_json::Value::Array(arr) => {
65 let values: Vec<Value> = arr.iter().map(|v| json_to_value(v, span)).collect();
66 Value::list(values, span)
67 }
68 serde_json::Value::Object(obj) => {
69 let mut record = Record::new();
70 for (k, v) in obj {
71 record.push(k, json_to_value(v, span));
72 }
73 Value::record(record, span)
74 }
75 }
76}
77
78pub fn frame_to_value(frame: &Frame, span: Span, with_timestamp: bool) -> Value {
79 let mut record = Record::new();
80
81 record.push("id", Value::string(frame.id.to_string(), span));
82 record.push("topic", Value::string(frame.topic.clone(), span));
83
84 if let Some(hash) = &frame.hash {
85 record.push("hash", Value::string(hash.to_string(), span));
86 }
87
88 if let Some(meta) = &frame.meta {
89 record.push("meta", json_to_value(meta, span));
90 }
91
92 if let Some(ttl) = &frame.ttl {
93 let ttl_str = match ttl {
94 crate::store::TTL::Forever => "forever".to_string(),
95 crate::store::TTL::Ephemeral => "ephemeral".to_string(),
96 crate::store::TTL::Time(duration) => format!("{}s", duration.as_secs()),
97 crate::store::TTL::Last(n) => format!("last:{n}"),
98 };
99 record.push("ttl", Value::string(ttl_str, span));
100 }
101
102 if with_timestamp {
103 let millis = frame.id.timestamp() as i64;
104 let dt: DateTime<Utc> = DateTime::from_timestamp_millis(millis)
105 .unwrap_or_else(|| DateTime::from_timestamp(0, 0).unwrap());
106 record.push("timestamp", Value::date(dt.fixed_offset(), span));
107 }
108
109 Value::record(record, span)
110}
111
112pub fn frame_to_pipeline(frame: &Frame, with_timestamp: bool) -> PipelineData {
113 PipelineData::Value(frame_to_value(frame, Span::unknown(), with_timestamp), None)
114}
115
116pub fn value_to_json(value: &Value) -> serde_json::Value {
117 match value {
118 Value::Nothing { .. } => serde_json::Value::Null,
119 Value::Bool { val, .. } => serde_json::Value::Bool(*val),
120 Value::Int { val, .. } => serde_json::Value::Number((*val).into()),
121 Value::Float { val, .. } => serde_json::Number::from_f64(*val)
122 .map(serde_json::Value::Number)
123 .unwrap_or(serde_json::Value::Null),
124 Value::String { val, .. } => serde_json::Value::String(val.clone()),
125 Value::Filesize { val, .. } => serde_json::Value::Number(val.get().into()),
126 Value::Date { val, .. } => serde_json::Value::String(val.to_rfc3339()),
127 Value::List { vals, .. } => {
128 serde_json::Value::Array(vals.iter().map(value_to_json).collect())
129 }
130 Value::Record { val, .. } => {
131 let mut map = serde_json::Map::new();
132 for (k, v) in val.iter() {
133 map.insert(k.clone(), value_to_json(v));
134 }
135 serde_json::Value::Object(map)
136 }
137 _ => serde_json::Value::Null,
138 }
139}
140
141pub fn write_pipeline_to_cas(
142 input: PipelineData,
143 store: &Store,
144 span: Span,
145) -> Result<Option<ssri::Integrity>, Box<ShellError>> {
146 let mut writer = store.cas_writer_sync().map_err(|e| {
147 Box::new(ShellError::Generic(GenericError::new(
148 "I/O Error",
149 e.to_string(),
150 span,
151 )))
152 })?;
153
154 match input {
155 PipelineData::Value(value, _) => match value {
156 Value::Nothing { .. } => Ok(None),
157 Value::String { val, .. } => {
158 writer.write_all(val.as_bytes()).map_err(|e| {
159 Box::new(ShellError::Generic(GenericError::new(
160 "I/O Error",
161 e.to_string(),
162 span,
163 )))
164 })?;
165
166 let hash = writer.commit().map_err(|e| {
167 Box::new(ShellError::Generic(GenericError::new(
168 "I/O Error",
169 e.to_string(),
170 span,
171 )))
172 })?;
173
174 Ok(Some(hash))
175 }
176 Value::Binary { val, .. } => {
177 writer.write_all(&val).map_err(|e| {
178 Box::new(ShellError::Generic(GenericError::new(
179 "I/O Error",
180 e.to_string(),
181 span,
182 )))
183 })?;
184
185 let hash = writer.commit().map_err(|e| {
186 Box::new(ShellError::Generic(GenericError::new(
187 "I/O Error",
188 e.to_string(),
189 span,
190 )))
191 })?;
192
193 Ok(Some(hash))
194 }
195 Value::Record { .. } => {
196 let json = value_to_json(&value);
197 let json_string = serde_json::to_string(&json).map_err(|e| {
198 Box::new(ShellError::Generic(GenericError::new(
199 "I/O Error",
200 e.to_string(),
201 span,
202 )))
203 })?;
204
205 writer.write_all(json_string.as_bytes()).map_err(|e| {
206 Box::new(ShellError::Generic(GenericError::new(
207 "I/O Error",
208 e.to_string(),
209 span,
210 )))
211 })?;
212
213 let hash = writer.commit().map_err(|e| {
214 Box::new(ShellError::Generic(GenericError::new(
215 "I/O Error",
216 e.to_string(),
217 span,
218 )))
219 })?;
220
221 Ok(Some(hash))
222 }
223 _ => Err(Box::new(ShellError::PipelineMismatch {
224 exp_input_type: format!(
225 "expected: string, binary, record, or nothing :: received: {typ:?}",
226 typ = value.get_type()
227 ),
228 dst_span: span,
229 src_span: value.span(),
230 })),
231 },
232 PipelineData::ListStream(_stream, ..) => {
233 panic!("ListStream handling is not yet implemented");
234 }
235 PipelineData::ByteStream(stream, ..) => {
236 if let Some(mut reader) = stream.reader() {
237 let mut buffer = [0; 8192];
238 loop {
239 let bytes_read = reader.read(&mut buffer).map_err(|e| {
240 Box::new(ShellError::Generic(GenericError::new(
241 "I/O Error",
242 e.to_string(),
243 span,
244 )))
245 })?;
246
247 if bytes_read == 0 {
248 break;
249 }
250
251 writer.write_all(&buffer[..bytes_read]).map_err(|e| {
252 Box::new(ShellError::Generic(GenericError::new(
253 "I/O Error",
254 e.to_string(),
255 span,
256 )))
257 })?;
258 }
259 }
260
261 let hash = writer.commit().map_err(|e| {
262 Box::new(ShellError::Generic(GenericError::new(
263 "I/O Error",
264 e.to_string(),
265 span,
266 )))
267 })?;
268
269 Ok(Some(hash))
270 }
271 PipelineData::Empty => Ok(None),
272 }
273}