nyx_space/od/msr/trackingdata/
io_ccsds_tdm.rs1use crate::io::ExportCfg;
20use crate::io::watermark::prj_name_ver;
21use crate::io::{InputOutputError, StdIOSnafu};
22use crate::od::msr::{Measurement, MeasurementType};
23use anise::constants::SPEED_OF_LIGHT_KM_S;
24use hifitime::efmt::{Format, Formatter};
25use hifitime::{Duration, Epoch, TimeScale};
26use indexmap::{IndexMap, IndexSet};
27use log::{error, info, warn};
28use snafu::ResultExt;
29use std::collections::HashMap;
30use std::fs::File;
31use std::io::Write;
32use std::io::{BufRead, BufReader, BufWriter};
33use std::path::{Path, PathBuf};
34use std::str::FromStr;
35
36use super::TrackingDataArc;
37
38impl TrackingDataArc {
39 pub fn from_tdm<P: AsRef<Path>>(
87 path: P,
88 aliases: Option<HashMap<String, String>>,
89 ) -> Result<Self, InputOutputError> {
90 let file = File::open(&path).context(StdIOSnafu {
91 action: "opening CCSDS TDM file for tracking arc",
92 })?;
93
94 let source = path.as_ref().to_path_buf().display().to_string();
95 info!("parsing CCSDS TDM {source}");
96
97 let mut measurements = Vec::new();
98 let mut metadata = HashMap::new();
99
100 let reader = BufReader::new(file);
101
102 let mut in_data_section = false;
103 let mut current_tracker = String::new();
104 let mut time_system = TimeScale::UTC;
105 let mut has_freq_data = false;
106 let mut msr_divider = 1.0;
107
108 for line in reader.lines() {
109 let line = line.context(StdIOSnafu {
110 action: "reading CCSDS TDM file",
111 })?;
112 let line = line.trim();
113
114 if line == "DATA_START" {
115 in_data_section = true;
116 continue;
117 } else if line == "DATA_STOP" {
118 in_data_section = false;
119 }
120
121 if !in_data_section {
122 if line.starts_with("PARTICIPANT_1") {
123 current_tracker = line.split('=').nth(1).unwrap_or("").trim().to_string();
124 if let Some(aliases) = &aliases
126 && let Some(alias) = aliases.get(¤t_tracker)
127 {
128 current_tracker = alias.clone();
129 }
130 } else if line.starts_with("TIME_SYSTEM") {
131 let ts = line.split('=').nth(1).unwrap_or("UTC").trim();
132 if let Ok(ts) = TimeScale::from_str(ts) {
134 time_system = ts;
135 } else {
136 return Err(InputOutputError::UnsupportedData {
137 which: format!("time scale `{ts}` not supported"),
138 });
139 }
140 } else if line.starts_with("PATH") {
141 match line.split(",").count() {
142 2 => msr_divider = 1.0,
143 3 => msr_divider = 2.0,
144 cnt => {
145 return Err(InputOutputError::UnsupportedData {
146 which: format!(
147 "found {cnt} paths in TDM, only 1 or 2 are supported"
148 ),
149 });
150 }
151 }
152 }
153
154 let mut splt = line.split('=');
155 if let Some(keyword) = splt.nth(0) {
156 if let Some(value) = splt.nth(0) {
158 metadata.insert(keyword.trim().to_string(), value.trim().to_string());
159 }
160 }
161
162 continue;
163 }
164
165 if let Some((mtype, epoch, value)) = parse_measurement_line(line, time_system)? {
166 let effective_divider = if [
169 MeasurementType::ReceiveFrequency,
170 MeasurementType::TransmitFrequency,
171 MeasurementType::TransmitFrequencyRate,
172 ]
173 .contains(&mtype)
174 {
175 has_freq_data = true;
176 1.0
177 } else {
178 msr_divider
179 };
180
181 let mut scaled_value = value;
182 if mtype == MeasurementType::Range {
183 if let Some(range_units) = metadata.get("RANGE_UNITS") {
184 match range_units.as_str() {
185 "km" => {
186 scaled_value /= effective_divider;
188 }
189 "RU" => {
190 return Err(InputOutputError::UnsupportedData {
191 which: "RANGE_UNITS `RU` requires mission-specific conversion and is not currently supported".to_string(),
192 });
193 }
194 "s" => {
195 scaled_value =
198 (scaled_value * SPEED_OF_LIGHT_KM_S) / effective_divider;
199 }
200 "m" => {
201 warn!(
202 "RANGE_UNITS in TDM file is `m`, which is not CCSDS compliant. Proceeding with conversion to km."
203 );
204 scaled_value = (scaled_value / 1000.0) / effective_divider;
205 }
206 "ms" => {
207 warn!(
208 "RANGE_UNITS in TDM file is `ms`, which is not CCSDS compliant. Proceeding with conversion to km."
209 );
210 scaled_value =
211 (scaled_value * 1e-3 * SPEED_OF_LIGHT_KM_S) / effective_divider;
212 }
213 "us" => {
214 warn!(
215 "RANGE_UNITS in TDM file is `us`, which is not CCSDS compliant. Proceeding with conversion to km."
216 );
217 scaled_value =
218 (scaled_value * 1e-6 * SPEED_OF_LIGHT_KM_S) / effective_divider;
219 }
220 "ns" => {
221 warn!(
222 "RANGE_UNITS in TDM file is `ns`, which is not CCSDS compliant. Proceeding with conversion to km."
223 );
224 scaled_value =
225 (scaled_value * 1e-9 * SPEED_OF_LIGHT_KM_S) / effective_divider;
226 }
227 _ => {
228 return Err(InputOutputError::UnsupportedData {
229 which: format!("unsupported RANGE_UNITS `{range_units}`"),
230 });
231 }
232 }
233 } else {
234 return Err(InputOutputError::MissingData {
235 which: "RANGE_UNITS not specified in metadata for RANGE measurement"
236 .to_string(),
237 });
238 }
239 } else {
240 scaled_value /= effective_divider;
241 }
242
243 let is_concurrent = measurements.last().is_some_and(|last: &Measurement| {
246 last.epoch == epoch && last.tracker == current_tracker
247 });
248
249 if is_concurrent {
250 measurements
251 .last_mut()
252 .unwrap() .data
254 .insert(mtype, scaled_value);
255 } else {
256 let mut data = IndexMap::new();
258 data.insert(mtype, scaled_value);
259
260 measurements.push(Measurement {
261 tracker: current_tracker.clone(),
262 epoch,
263 data,
264 rejected: false,
265 });
266 }
267 }
268 }
269
270 let mut turnaround_ratio = None;
271 let drop_freq_data;
272 if has_freq_data {
273 if let Some(ta_num_str) = metadata.get("TURNAROUND_NUMERATOR") {
275 if let Some(ta_denom_str) = metadata.get("TURNAROUND_DENOMINATOR") {
276 if let Ok(ta_num) = ta_num_str.parse::<i32>() {
277 if let Ok(ta_denom) = ta_denom_str.parse::<i32>() {
278 turnaround_ratio = Some(f64::from(ta_num) / f64::from(ta_denom));
280 info!("turn-around ratio is {ta_num}/{ta_denom}");
281 drop_freq_data = false;
282 } else {
283 error!(
284 "turn-around denominator `{ta_denom_str}` is not a valid integer"
285 );
286 drop_freq_data = true;
287 }
288 } else {
289 error!("turn-around numerator `{ta_num_str}` is not a valid integer");
290 drop_freq_data = true;
291 }
292 } else {
293 error!(
294 "required turn-around denominator missing from metadata -- dropping ALL RECEIVE/TRANSMIT data"
295 );
296 drop_freq_data = true;
297 }
298 } else {
299 error!(
300 "required turn-around numerator missing from metadata -- dropping ALL RECEIVE/TRANSMIT data"
301 );
302 drop_freq_data = true;
303 }
304 } else {
305 drop_freq_data = true;
306 }
307
308 let corrections_applied = if let Some(corr_flag) = metadata.get("CORRECTIONS_APPLIED") {
309 match corr_flag.trim().to_lowercase().as_str() {
310 "no" => false,
311 "yes" => true,
312 _ => {
313 warn!("invalid CORRECTIONS_APPLIED `{corr_flag}`");
314 true
315 }
316 }
317 } else {
318 true
319 };
320
321 let mut freq_types = IndexSet::new();
324 freq_types.insert(MeasurementType::ReceiveFrequency);
325 freq_types.insert(MeasurementType::TransmitFrequency);
326 freq_types.insert(MeasurementType::TransmitFrequencyRate);
327
328 let mut latest_transmit_freq = None;
329 let mut latest_transmit_epoch = None;
330 let mut latest_transmit_rate = 0.0;
331
332 let mut all_applied_corrections = IndexSet::new();
333
334 for measurement in &mut measurements {
335 let epoch = measurement.epoch;
336 if !corrections_applied {
338 for msr_type in [
339 MeasurementType::Range,
340 MeasurementType::Doppler,
341 MeasurementType::Azimuth,
342 MeasurementType::Elevation,
343 MeasurementType::ReceiveFrequency,
344 MeasurementType::TransmitFrequency,
345 MeasurementType::TransmitFrequencyRate,
346 ] {
347 let kw = format!("CORRECTION_{}", msr_type.ccsds_tdm_name());
348 if let Some(correction_str) = metadata.get(&kw) {
349 if let Ok(correction) = correction_str.parse::<f64>() {
350 measurement.correct(msr_type, correction);
351 all_applied_corrections.insert(msr_type);
352 } else {
353 warn!("invalid correction value for {kw}");
354 }
355 }
356 }
357 }
358
359 if drop_freq_data {
360 for freq in &freq_types {
361 measurement.data.swap_remove(freq);
362 }
363 continue;
364 }
365
366 if let Some(rate) = measurement
368 .data
369 .get(&MeasurementType::TransmitFrequencyRate)
370 {
371 if let (Some(last_f), Some(last_e)) = (latest_transmit_freq, latest_transmit_epoch)
372 {
373 let dt: Duration = epoch - last_e;
374 latest_transmit_freq = Some(last_f + latest_transmit_rate * dt.to_seconds());
375 }
376 latest_transmit_epoch = Some(epoch);
377 latest_transmit_rate = *rate;
378 }
379
380 if let Some(freq) = measurement.data.get(&MeasurementType::TransmitFrequency) {
381 latest_transmit_freq = Some(*freq);
382 latest_transmit_epoch = Some(epoch);
383 }
384
385 if !measurement
386 .data
387 .contains_key(&MeasurementType::ReceiveFrequency)
388 {
389 for freq in &freq_types {
392 measurement.data.swap_remove(freq);
393 }
394 continue;
395 }
396
397 if latest_transmit_freq.is_none() {
399 warn!(
400 "receive frequency found at {epoch} but no transmit frequency was ever set, ignoring"
401 );
402 for freq in &freq_types {
403 measurement.data.swap_remove(freq);
404 }
405 continue;
406 }
407
408 let dt: Duration = epoch - latest_transmit_epoch.unwrap();
409 let transmit_freq_hz =
410 latest_transmit_freq.unwrap() + latest_transmit_rate * dt.to_seconds();
411
412 let receive_freq_hz = *measurement
413 .data
414 .get(&MeasurementType::ReceiveFrequency)
415 .unwrap();
416
417 let doppler_shift_hz = transmit_freq_hz * turnaround_ratio.unwrap() - receive_freq_hz;
419 let rho_dot_km_s = (doppler_shift_hz * SPEED_OF_LIGHT_KM_S)
421 / (2.0 * transmit_freq_hz * turnaround_ratio.unwrap());
422
423 for freq in &freq_types {
425 measurement.data.swap_remove(freq);
426 }
427 measurement
428 .data
429 .insert(MeasurementType::Doppler, rho_dot_km_s);
430 }
431
432 if !all_applied_corrections.is_empty() {
433 info!("applied corrections for {all_applied_corrections:?}");
434 }
435
436 let moduli = if let Some(range_modulus) = metadata.get("RANGE_MODULUS") {
437 if let Ok(value) = range_modulus.parse::<f64>() {
438 if value > 0.0 {
439 let mut modulos = IndexMap::new();
440 modulos.insert(MeasurementType::Range, value);
441 Some(modulos)
443 } else {
444 None
446 }
447 } else {
448 warn!("could not parse RANGE_MODULUS of `{range_modulus}` as a double");
449 None
450 }
451 } else {
452 None
453 };
454
455 measurements.retain(|m| !m.data.is_empty());
457
458 let mut trk = Self {
459 measurements,
460 source: Some(source),
461 moduli,
462 force_reject: false,
463 };
464
465 trk.sort();
467
468 if trk.unique_types().is_empty() {
469 Err(InputOutputError::EmptyDataset {
470 action: "CCSDS TDM file",
471 })
472 } else {
473 Ok(trk)
474 }
475 }
476
477 pub fn to_tdm_file<P: AsRef<Path>>(
479 mut self,
480 spacecraft_name: String,
481 aliases: Option<HashMap<String, String>>,
482 path: P,
483 cfg: ExportCfg,
484 ) -> Result<PathBuf, InputOutputError> {
485 if self.is_empty() {
486 return Err(InputOutputError::MissingData {
487 which: " - empty tracking data cannot be exported to TDM".to_string(),
488 });
489 }
490
491 if let Some(start_epoch) = cfg.start_epoch {
493 if let Some(end_epoch) = cfg.end_epoch {
494 self = self.filter_by_epoch(start_epoch..end_epoch);
495 } else {
496 self = self.filter_by_epoch(start_epoch..);
497 }
498 } else if let Some(end_epoch) = cfg.end_epoch {
499 self = self.filter_by_epoch(..end_epoch);
500 }
501
502 let tick = Epoch::now().unwrap();
503 info!("Exporting tracking data to CCSDS TDM file...");
504
505 let path_buf = cfg.actual_path(path);
507
508 let metadata = cfg.metadata.unwrap_or_default();
509
510 let file = File::create(&path_buf).context(StdIOSnafu {
511 action: "creating CCSDS TDM file for tracking arc",
512 })?;
513 let mut writer = BufWriter::new(file);
514
515 let err_hdlr = |source| InputOutputError::StdIOError {
516 source,
517 action: "writing data to TDM file",
518 };
519
520 let iso8601_no_ts = Format::from_str("%Y-%m-%dT%H:%M:%S.%f").unwrap();
522
523 writeln!(writer, "CCSDS_TDM_VERS = 2.0").map_err(err_hdlr)?;
525 writeln!(
526 writer,
527 "\nCOMMENT Build by {} -- https://nyxspace.com",
528 prj_name_ver()
529 )
530 .map_err(err_hdlr)?;
531 writeln!(
532 writer,
533 "COMMENT Nyx Space provided under the AGPL v3 open source license -- https://nyxspace.com/pricing\n"
534 )
535 .map_err(err_hdlr)?;
536 writeln!(
537 writer,
538 "CREATION_DATE = {}",
539 Formatter::new(Epoch::now().unwrap(), iso8601_no_ts)
540 )
541 .map_err(err_hdlr)?;
542 writeln!(
543 writer,
544 "ORIGINATOR = {}\n",
545 metadata
546 .get("originator")
547 .unwrap_or(&"Nyx Space".to_string())
548 )
549 .map_err(err_hdlr)?;
550
551 let trackers = self.unique_aliases();
554
555 for tracker in trackers {
556 let tracker_data = self.clone().filter_by_tracker(tracker.clone());
557
558 let types = tracker_data.unique_types();
559
560 let two_way_types = types
561 .iter()
562 .filter(|msr_type| msr_type.may_be_two_way())
563 .copied()
564 .collect::<Vec<_>>();
565
566 let one_way_types = types
567 .iter()
568 .filter(|msr_type| !msr_type.may_be_two_way())
569 .copied()
570 .collect::<Vec<_>>();
571
572 for (tno, types) in [two_way_types, one_way_types].iter().enumerate() {
574 writeln!(writer, "META_START").map_err(err_hdlr)?;
575 writeln!(writer, "\tTIME_SYSTEM = UTC").map_err(err_hdlr)?;
576 writeln!(
577 writer,
578 "\tSTART_TIME = {}",
579 Formatter::new(tracker_data.start_epoch().unwrap(), iso8601_no_ts)
580 )
581 .map_err(err_hdlr)?;
582 writeln!(
583 writer,
584 "\tSTOP_TIME = {}",
585 Formatter::new(tracker_data.end_epoch().unwrap(), iso8601_no_ts)
586 )
587 .map_err(err_hdlr)?;
588
589 let multiplier = if tno == 0 {
590 writeln!(writer, "\tPATH = 1,2,1").map_err(err_hdlr)?;
591 2.0
592 } else {
593 writeln!(writer, "\tPATH = 1,2").map_err(err_hdlr)?;
594 1.0
595 };
596
597 writeln!(
598 writer,
599 "\tPARTICIPANT_1 = {}",
600 if let Some(aliases) = &aliases {
601 if let Some(alias) = aliases.get(&tracker) {
602 alias
603 } else {
604 &tracker
605 }
606 } else {
607 &tracker
608 }
609 )
610 .map_err(err_hdlr)?;
611
612 writeln!(writer, "\tPARTICIPANT_2 = {spacecraft_name}").map_err(err_hdlr)?;
613
614 writeln!(writer, "\tMODE = SEQUENTIAL").map_err(err_hdlr)?;
615
616 for (k, v) in &metadata {
618 if k != "originator" {
619 writeln!(writer, "\t{k} = {v}").map_err(err_hdlr)?;
620 }
621 }
622
623 if types.contains(&MeasurementType::Range) {
624 writeln!(writer, "\tRANGE_UNITS = km").map_err(err_hdlr)?;
625
626 if let Some(moduli) = &self.moduli
627 && let Some(range_modulus) = moduli.get(&MeasurementType::Range)
628 {
629 writeln!(writer, "\tRANGE_MODULUS = {range_modulus:E}")
630 .map_err(err_hdlr)?;
631 }
632 }
633
634 if types.contains(&MeasurementType::Azimuth)
635 || types.contains(&MeasurementType::Elevation)
636 {
637 writeln!(writer, "\tANGLE_TYPE = AZEL").map_err(err_hdlr)?;
638 }
639
640 writeln!(writer, "META_STOP\n").map_err(err_hdlr)?;
641
642 writeln!(writer, "DATA_START").map_err(err_hdlr)?;
644
645 for m in &tracker_data.measurements {
647 for (mtype, value) in &m.data {
648 if !types.contains(mtype) {
649 continue;
650 }
651
652 writeln!(
653 writer,
654 "\t{:<20} = {:<23}\t{:.12}",
655 mtype.ccsds_tdm_name(),
656 Formatter::new(m.epoch, iso8601_no_ts),
657 value * multiplier
658 )
659 .map_err(err_hdlr)?;
660 }
661 }
662
663 writeln!(writer, "DATA_STOP\n").map_err(err_hdlr)?;
664 }
665 }
666
667 #[allow(clippy::writeln_empty_string)]
668 writeln!(writer, "").map_err(err_hdlr)?;
669
670 let tock_time = Epoch::now().unwrap() - tick;
672 info!("CCSDS TDM written to {} in {tock_time}", path_buf.display());
673 Ok(path_buf)
674 }
675}
676
677fn parse_measurement_line(
678 line: &str,
679 time_system: TimeScale,
680) -> Result<Option<(MeasurementType, Epoch, f64)>, InputOutputError> {
681 let parts: Vec<&str> = line.split('=').collect();
682 if parts.len() != 2 {
683 return Ok(None);
684 }
685
686 let (mtype_str, data) = (parts[0].trim(), parts[1].trim());
687 let mtype = match mtype_str {
688 "RANGE" => MeasurementType::Range,
689 "DOPPLER_INSTANTANEOUS" | "DOPPLER_INTEGRATED" => MeasurementType::Doppler,
690 "ANGLE_1" => MeasurementType::Azimuth,
691 "ANGLE_2" => MeasurementType::Elevation,
692 "RECEIVE_FREQ" | "RECEIVE_FREQ_1" | "RECEIVE_FREQ_2" | "RECEIVE_FREQ_3"
693 | "RECEIVE_FREQ_4" | "RECEIVE_FREQ_5" => MeasurementType::ReceiveFrequency,
694 "TRANSMIT_FREQ" | "TRANSMIT_FREQ_1" | "TRANSMIT_FREQ_2" | "TRANSMIT_FREQ_3"
695 | "TRANSMIT_FREQ_4" | "TRANSMIT_FREQ_5" => MeasurementType::TransmitFrequency,
696 "TRANSMIT_FREQ_RATE"
697 | "TRANSMIT_FREQ_RATE_1"
698 | "TRANSMIT_FREQ_RATE_2"
699 | "TRANSMIT_FREQ_RATE_3"
700 | "TRANSMIT_FREQ_RATE_4"
701 | "TRANSMIT_FREQ_RATE_5" => MeasurementType::TransmitFrequencyRate,
702 _ => {
703 return Err(InputOutputError::UnsupportedData {
704 which: mtype_str.to_string(),
705 });
706 }
707 };
708
709 let data_parts: Vec<&str> = data.split_whitespace().collect();
710 if data_parts.len() != 2 {
711 return Ok(None);
712 }
713
714 let epoch =
715 Epoch::from_gregorian_str(&format!("{} {time_system}", data_parts[0])).map_err(|e| {
716 InputOutputError::Inconsistency {
717 msg: format!("{e} when parsing epoch"),
718 }
719 })?;
720
721 let value = data_parts[1]
722 .parse::<f64>()
723 .map_err(|e| InputOutputError::UnsupportedData {
724 which: format!("`{}` is not a float: {e}", data_parts[1]),
725 })?;
726
727 Ok(Some((mtype, epoch, value)))
728}