Skip to main content

fits_io/
slice_bin_table_hdu.rs

1use crate::bin_table::BinTable;
2#[cfg(feature = "tokio")]
3use crate::bin_table::{FieldDefinition, OwnedRow};
4use crate::hdu::{BinTableHDU, HDU};
5use crate::header::Header;
6#[cfg(feature = "tokio")]
7use futures::StreamExt;
8#[cfg(feature = "tokio")]
9use futures::stream;
10#[cfg(feature = "tokio")]
11use futures::stream::BoxStream;
12#[cfg(feature = "serde")]
13use serde::de::DeserializeOwned;
14use std::error::Error;
15use std::sync::Arc;
16
17/// A binary table HDU backed by a buffer rather than a file.
18#[derive(Debug, Clone)]
19pub struct SliceBinTableHDU {
20    header: Header,
21    data: Arc<Vec<u8>>,
22    data_offset: usize,
23    /// A table set through [`BinTableHDU::set_table`], which stands in for
24    /// whatever the buffer holds.
25    pending: Option<Vec<u8>>,
26}
27
28impl SliceBinTableHDU {
29    pub(crate) fn new(header: Header, data: Arc<Vec<u8>>, data_offset: usize) -> Self {
30        Self {
31            header,
32            data,
33            data_offset,
34            pending: None,
35        }
36    }
37
38    /// An HDU holding `table` and nothing read from a buffer.
39    pub fn from_table(table: &BinTable) -> Result<Self, Box<dyn Error + Send + Sync>> {
40        let mut hdu = Self {
41            header: Header::default(),
42            data: Arc::new(Vec::new()),
43            data_offset: 0,
44            pending: None,
45        };
46        hdu.set_table(table)?;
47
48        Ok(hdu)
49    }
50
51    /// This HDU's whole data section: the rows, then the heap.
52    pub(crate) fn data_bytes(&self) -> &[u8] {
53        if let Some(pending) = &self.pending {
54            return pending;
55        }
56
57        let len = self.header.data_bytes_len();
58        self.data
59            .get(self.data_offset..)
60            .and_then(|data| data.get(..len))
61            .unwrap_or_default()
62    }
63}
64
65impl HDU for SliceBinTableHDU {
66    fn header(&self) -> &Header {
67        &self.header
68    }
69
70    fn header_mut(&mut self) -> &mut Header {
71        &mut self.header
72    }
73}
74
75impl BinTableHDU for SliceBinTableHDU {
76    fn table_data_bytes_len(&self) -> u64 {
77        let (Some(bytes_per_row), Some(rows)) = (self.header.naxis_n(0), self.header.naxis_n(1))
78        else {
79            return 0;
80        };
81
82        (bytes_per_row.max(0) as u64).saturating_mul(rows.max(0) as u64)
83    }
84
85    fn read_table(&self) -> Result<BinTable, Box<dyn Error + Send + Sync>> {
86        BinTable::from_u8(&self.header, self.data_bytes().to_vec())
87    }
88
89    fn set_table(&mut self, table: &BinTable) -> Result<(), Box<dyn Error + Send + Sync>> {
90        crate::bin_table::write::apply_to_header(table, &mut self.header);
91        self.pending = Some(table.data().to_vec());
92
93        Ok(())
94    }
95
96    #[cfg(feature = "serde")]
97    fn read_rows<T: DeserializeOwned + Send + Sync>(
98        &self,
99    ) -> Result<Vec<T>, Box<dyn Error + Send + Sync>> {
100        let table = self.read_table()?;
101        Ok(crate::bin_table::from_bin_table(&table)?)
102    }
103
104    #[cfg(feature = "tokio")]
105    fn stream_table_rows_raw(
106        &self,
107    ) -> Result<BoxStream<'_, OwnedRow>, Box<dyn Error + Send + Sync>> {
108        let bytes_per_row = self.header.naxis_n(0).unwrap_or(0);
109        let bytes_per_row = usize::try_from(bytes_per_row)
110            .map_err(|_| format!("NAXIS1 must not be negative, but was {}", bytes_per_row))?;
111
112        if bytes_per_row == 0 {
113            return Ok(stream::empty().boxed());
114        }
115
116        let field_definitions = Arc::new(FieldDefinition::all_from_header(&self.header)?);
117
118        let data = self.data_bytes();
119        let heap_offset = self.header.table_heap_offset().min(data.len());
120        let heap = Arc::new(data[heap_offset..].to_vec());
121
122        // Everything is already in memory, so the stream exists to match the
123        // file-backed API rather than to keep memory down.
124        let rows: Vec<_> = data[..heap_offset]
125            .chunks_exact(bytes_per_row)
126            .map(|row| {
127                OwnedRow::new(
128                    Arc::clone(&field_definitions),
129                    Arc::clone(&heap),
130                    row.to_vec(),
131                )
132            })
133            .collect();
134
135        Ok(stream::iter(rows).boxed())
136    }
137
138    #[cfg(feature = "serde")]
139    #[cfg(feature = "tokio")]
140    fn stream_table_rows<T: DeserializeOwned + Send + Sync>(
141        &self,
142    ) -> Result<BoxStream<'_, crate::Result<T>>, Box<dyn Error + Send + Sync>> {
143        Ok(self
144            .stream_table_rows_raw()?
145            .map(|row| crate::bin_table::from_bin_table_row(&row.row()))
146            .boxed())
147    }
148}