Skip to content

Commit 58c0ea2

Browse files
committed
read implement in progress
1 parent aa64786 commit 58c0ea2

9 files changed

Lines changed: 275 additions & 22 deletions

File tree

src/atomic_message/mod.rs

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,24 @@
1+
mod read;
2+
3+
use std::sync::Arc;
4+
5+
use eccodes_sys::codes_handle;
6+
7+
use crate::{CodesHandle, codes_handle::ThreadSafeHandle};
8+
9+
pub use read::{KeyRead, ArrayKeyRead, ScalarKeyRead};
10+
11+
/// Because standard `KeyedMessage` is not Copy or Clone it can provide access methods without
12+
/// requiring `&mut self`. As `AtomicMessage` implements `Send + Sync` this exclusive method access is not
13+
/// guaranteed with just `&self`. `AtomicMessage` also implements a minimal subset of functionalities
14+
/// to limit the risk of some internal ecCodes functions not being thread-safe.
15+
///
16+
/// Right now `AtomicMessage` is also not clonable
17+
#[derive(Debug)]
18+
pub struct AtomicMessage<S: ThreadSafeHandle> {
19+
pub(crate) _parent: Arc<CodesHandle<S>>,
20+
pub(crate) message_handle: *mut codes_handle,
21+
}
22+
23+
unsafe impl<S: ThreadSafeHandle> Send for AtomicMessage<S> {}
24+
unsafe impl<S: ThreadSafeHandle> Sync for AtomicMessage<S> {}

src/atomic_message/read.rs

Lines changed: 110 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
1+
use std::cmp::Ordering;
2+
3+
use crate::{
4+
CodesError,
5+
atomic_message::AtomicMessage,
6+
codes_handle::ThreadSafeHandle,
7+
intermediate_bindings::{
8+
NativeKeyType, codes_get_bytes, codes_get_double, codes_get_double_array, codes_get_long,
9+
codes_get_long_array, codes_get_native_type, codes_get_size, codes_get_string,
10+
},
11+
};
12+
13+
#[doc(hidden)]
14+
pub trait KeyReadHelpers {
15+
fn get_key_size(&mut self, key_name: &str) -> Result<usize, CodesError>;
16+
fn get_key_native_type(&mut self, key_name: &str) -> Result<NativeKeyType, CodesError>;
17+
}
18+
pub trait KeyRead<T>: KeyReadHelpers {
19+
fn read_key_unchecked(&mut self, name: &str) -> Result<T, CodesError>;
20+
}
21+
22+
pub trait ArrayKeyRead<T>: KeyRead<T> {
23+
fn read_key(&mut self, key_name: &str) -> Result<T, CodesError> {
24+
match self.get_key_native_type(key_name)? {
25+
NativeKeyType::Bytes => (),
26+
_ => return Err(CodesError::WrongRequestedKeyType),
27+
}
28+
29+
let key_size = self.get_key_size(key_name)?;
30+
31+
if key_size < 1 {
32+
return Err(CodesError::IncorrectKeySize);
33+
}
34+
35+
self.read_key_unchecked(key_name)
36+
}
37+
}
38+
39+
pub trait ScalarKeyRead<T>: KeyRead<T> {
40+
fn read_key(&mut self, key_name: &str) -> Result<T, CodesError> {
41+
match self.get_key_native_type(key_name)? {
42+
NativeKeyType::Long => (),
43+
_ => return Err(CodesError::WrongRequestedKeyType),
44+
}
45+
46+
let key_size = self.get_key_size(key_name)?;
47+
48+
match key_size.cmp(&1) {
49+
Ordering::Greater => return Err(CodesError::WrongRequestedKeySize),
50+
Ordering::Less => return Err(CodesError::IncorrectKeySize),
51+
Ordering::Equal => (),
52+
}
53+
54+
self.read_key_unchecked(key_name)
55+
}
56+
}
57+
58+
impl<S: ThreadSafeHandle> KeyReadHelpers for AtomicMessage<S> {
59+
fn get_key_size(&mut self, key_name: &str) -> Result<usize, CodesError> {
60+
unsafe { codes_get_size(self.message_handle, key_name) }
61+
}
62+
63+
fn get_key_native_type(&mut self, key_name: &str) -> Result<NativeKeyType, CodesError> {
64+
unsafe { codes_get_native_type(self.message_handle, key_name) }
65+
}
66+
}
67+
68+
impl<S: ThreadSafeHandle> KeyRead<i64> for AtomicMessage<S> {
69+
fn read_key_unchecked(&mut self, key_name: &str) -> Result<i64, CodesError> {
70+
unsafe { codes_get_long(self.message_handle, key_name) }
71+
}
72+
}
73+
74+
impl<S: ThreadSafeHandle> KeyRead<f64> for AtomicMessage<S> {
75+
fn read_key_unchecked(&mut self, key_name: &str) -> Result<f64, CodesError> {
76+
unsafe { codes_get_double(self.message_handle, key_name) }
77+
}
78+
}
79+
80+
impl<S: ThreadSafeHandle> KeyRead<String> for AtomicMessage<S> {
81+
fn read_key_unchecked(&mut self, key_name: &str) -> Result<String, CodesError> {
82+
unsafe { codes_get_string(self.message_handle, key_name) }
83+
}
84+
}
85+
86+
impl<S: ThreadSafeHandle> KeyRead<Vec<i64>> for AtomicMessage<S> {
87+
fn read_key_unchecked(&mut self, key_name: &str) -> Result<Vec<i64>, CodesError> {
88+
unsafe { codes_get_long_array(self.message_handle, key_name) }
89+
}
90+
}
91+
92+
impl<S: ThreadSafeHandle> KeyRead<Vec<f64>> for AtomicMessage<S> {
93+
fn read_key_unchecked(&mut self, key_name: &str) -> Result<Vec<f64>, CodesError> {
94+
unsafe { codes_get_double_array(self.message_handle, key_name) }
95+
}
96+
}
97+
98+
impl<S: ThreadSafeHandle> KeyRead<Vec<u8>> for AtomicMessage<S> {
99+
fn read_key_unchecked(&mut self, key_name: &str) -> Result<Vec<u8>, CodesError> {
100+
unsafe { codes_get_bytes(self.message_handle, key_name) }
101+
}
102+
}
103+
104+
impl<S: ThreadSafeHandle> ScalarKeyRead<i64> for AtomicMessage<S> {}
105+
impl<S: ThreadSafeHandle> ScalarKeyRead<f64> for AtomicMessage<S> {}
106+
107+
impl<S: ThreadSafeHandle> ArrayKeyRead<String> for AtomicMessage<S> {}
108+
impl<S: ThreadSafeHandle> ArrayKeyRead<Vec<f64>> for AtomicMessage<S> {}
109+
impl<S: ThreadSafeHandle> ArrayKeyRead<Vec<i64>> for AtomicMessage<S> {}
110+
impl<S: ThreadSafeHandle> ArrayKeyRead<Vec<u8>> for AtomicMessage<S> {}

src/codes_handle/atomic_iterator.rs

Lines changed: 5 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -1,11 +1,10 @@
1-
#![allow(unused)]
2-
31
use std::sync::Arc;
42

5-
use eccodes_sys::codes_handle;
63
use fallible_iterator::FallibleIterator;
74

8-
use crate::{CodesError, CodesHandle, codes_handle::ThreadSafeHandle};
5+
use crate::{
6+
CodesError, CodesHandle, atomic_message::AtomicMessage, codes_handle::ThreadSafeHandle,
7+
};
98

109
#[derive(Debug)]
1110
pub struct AtomicMessageGenerator<S: ThreadSafeHandle> {
@@ -32,20 +31,12 @@ impl<S: ThreadSafeHandle> FallibleIterator for AtomicMessageGenerator<S> {
3231
} else {
3332
Ok(Some(AtomicMessage {
3433
_parent: self.codes_handle.clone(),
35-
pointer: new_eccodes_handle,
34+
message_handle: new_eccodes_handle,
3635
}))
3736
}
3837
}
3938
}
4039

41-
#[derive(Debug)]
42-
pub struct AtomicMessage<S: ThreadSafeHandle> {
43-
_parent: Arc<CodesHandle<S>>,
44-
pointer: *mut codes_handle,
45-
}
46-
unsafe impl<S: ThreadSafeHandle> Send for AtomicMessage<S> {}
47-
unsafe impl<S: ThreadSafeHandle> Sync for AtomicMessage<S> {}
48-
4940
#[cfg(test)]
5041
mod tests {
5142
use std::{
@@ -79,7 +70,7 @@ mod tests {
7970
for _ in 0..1000 {
8071
b.wait();
8172
let _ = unsafe {
82-
crate::intermediate_bindings::codes_get_size(msg.pointer, "shortName")
73+
crate::intermediate_bindings::codes_get_size(msg.message_handle, "shortName")
8374
.unwrap()
8475
};
8576
}

src/codes_handle/iterator.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -39,7 +39,7 @@ impl<'ch, S: HandleGenerator> FallibleIterator for KeyedMessageGenerator<'ch, S>
3939
#[cfg(test)]
4040
mod tests {
4141
use crate::{
42-
DynamicKeyType, FallibleIterator,
42+
keyed_message::DynamicKeyType, FallibleIterator,
4343
codes_handle::{CodesHandle, ProductKind},
4444
};
4545
use anyhow::{Context, Ok, Result};

src/keyed_message/read.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
use std::cmp::Ordering;
22

33
use crate::{
4-
DynamicKeyType, KeyRead, KeyedMessage,
4+
keyed_message::DynamicKeyType, keyed_message::KeyRead, KeyedMessage,
55
errors::CodesError,
66
intermediate_bindings::{
77
NativeKeyType, codes_get_bytes, codes_get_double, codes_get_double_array, codes_get_long,
@@ -300,7 +300,7 @@ mod tests {
300300
use anyhow::{Context, Result};
301301

302302
use crate::codes_handle::{CodesHandle, ProductKind};
303-
use crate::{DynamicKeyType, FallibleIterator};
303+
use crate::{keyed_message::DynamicKeyType, FallibleIterator};
304304
use std::path::Path;
305305

306306
#[test]

src/keyed_message/write.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -135,7 +135,7 @@ mod tests {
135135
use fallible_iterator::FallibleIterator;
136136

137137
use crate::{
138-
DynamicKeyType, KeyWrite,
138+
keyed_message::DynamicKeyType, keyed_message::KeyWrite,
139139
codes_handle::{CodesHandle, ProductKind},
140140
};
141141
use std::{fs::remove_file, path::Path};

src/lib.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -222,6 +222,7 @@ pub mod keys_iterator;
222222
#[cfg_attr(docsrs, doc(cfg(feature = "message_ndarray")))]
223223
pub mod message_ndarray;
224224
mod pointer_guard;
225+
pub mod atomic_message;
225226

226227
pub use codes_handle::{CodesHandle, ProductKind};
227228
#[cfg(feature = "experimental_index")]
@@ -230,5 +231,6 @@ pub use codes_index::CodesIndex;
230231
pub use codes_nearest::{CodesNearest, NearestGridpoint};
231232
pub use errors::CodesError;
232233
pub use fallible_iterator::{FallibleIterator, IntoFallibleIterator};
233-
pub use keyed_message::{DynamicKeyType, KeyRead, KeyWrite, KeyedMessage};
234+
pub use keyed_message::{KeyedMessage};
234235
pub use keys_iterator::{KeysIterator, KeysIteratorFlags};
236+
pub use atomic_message::{AtomicMessage};

src/message_ndarray.rs

Lines changed: 128 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,13 @@
33
44
use ndarray::{Array2, Array3, s};
55

6-
use crate::{CodesError, KeyRead, KeyedMessage, errors::MessageNdarrayError};
6+
use crate::{
7+
AtomicMessage, CodesError, KeyedMessage,
8+
atomic_message::{ArrayKeyRead, ScalarKeyRead},
9+
codes_handle::ThreadSafeHandle,
10+
errors::MessageNdarrayError,
11+
keyed_message::KeyRead,
12+
};
713

814
/// Struct returned by [`KeyedMessage::to_lons_lats_values()`] method.
915
/// The arrays are collocated, meaning that `longitudes[i, j]` and `latitudes[i, j]` are the coordinates of `values[i, j]`.
@@ -138,15 +144,135 @@ impl KeyedMessage<'_> {
138144
}
139145
}
140146

147+
impl<S: ThreadSafeHandle> AtomicMessage<S> {
148+
/// Converts the message to a 2D ndarray.
149+
///
150+
/// Returns ndarray where first dimension represents y coordinates and second dimension represents x coordinates,
151+
/// ie. `[lat, lon]`.
152+
///
153+
/// Common convention for grib files on regular lon-lat grid assumes that:
154+
/// index `[0, 0]` is the top-left corner of the grid:
155+
/// x coordinates are increasing with the i index,
156+
/// y coordinates are decreasing with the j index.
157+
///
158+
/// This convention can be checked with `iScansNegatively` and `jScansPositively` keys -
159+
/// if both are false, the above convention is used.
160+
///
161+
/// Requires the keys `Ni`, `Nj` and `values` to be present in the message.
162+
///
163+
/// Tested only with simple lat-lon grids.
164+
///
165+
/// # Errors
166+
///
167+
/// - When the required keys are not present or if their values are not of the expected type
168+
/// - When the number of values mismatch with the `Ni` and `Nj` keys
169+
#[cfg_attr(docsrs, doc(cfg(feature = "message_ndarray")))]
170+
pub fn to_ndarray(&self) -> Result<Array2<f64>, CodesError> {
171+
let ni: i64 = self.read_key("Ni")?;
172+
let ni = usize::try_from(ni).map_err(MessageNdarrayError::from)?;
173+
174+
let nj: i64 = self.read_key("Nj")?;
175+
let nj = usize::try_from(nj).map_err(MessageNdarrayError::from)?;
176+
177+
let vals: Vec<f64> = self.read_key("values")?;
178+
if vals.len() != (ni * nj) {
179+
return Err(MessageNdarrayError::UnexpectedValuesLength(vals.len(), ni * nj).into());
180+
}
181+
182+
let j_scanning: i64 = self.read_key("jPointsAreConsecutive")?;
183+
184+
if ![0, 1].contains(&j_scanning) {
185+
return Err(MessageNdarrayError::UnexpectedKeyValue(
186+
"jPointsAreConsecutive".to_owned(),
187+
)
188+
.into());
189+
}
190+
191+
let j_scanning = j_scanning != 0;
192+
193+
let shape = if j_scanning { (ni, nj) } else { (nj, ni) };
194+
let vals = Array2::from_shape_vec(shape, vals).map_err(MessageNdarrayError::from)?;
195+
196+
if j_scanning {
197+
Ok(vals.reversed_axes())
198+
} else {
199+
Ok(vals)
200+
}
201+
}
202+
203+
/// Same as [`KeyedMessage::to_ndarray()`] but returns the longitudes and latitudes alongside values.
204+
/// Fields are returned as separate arrays in [`RustyCodesMessage`].
205+
///
206+
/// Compared to `to_ndarray` this method has performance overhead as returned arrays may be cloned.
207+
///
208+
/// This method requires the `latLonValues`, `Ni` and `Nj` keys to be present in the message.
209+
///
210+
/// # Errors
211+
///
212+
/// - When the required keys are not present or if their values are not of the expected type
213+
/// - When the number of values mismatch with the `Ni` and `Nj` keys
214+
#[cfg_attr(docsrs, doc(cfg(feature = "message_ndarray")))]
215+
pub fn to_lons_lats_values(&self) -> Result<RustyCodesMessage, CodesError> {
216+
let ni: i64 = self.read_key("Ni")?;
217+
let ni = usize::try_from(ni).map_err(MessageNdarrayError::from)?;
218+
219+
let nj: i64 = self.read_key("Nj")?;
220+
let nj = usize::try_from(nj).map_err(MessageNdarrayError::from)?;
221+
222+
let latlonvals: Vec<f64> = self.read_key("latLonValues")?;
223+
224+
if latlonvals.len() != (ni * nj * 3) {
225+
return Err(
226+
MessageNdarrayError::UnexpectedValuesLength(latlonvals.len(), ni * nj * 3).into(),
227+
);
228+
}
229+
230+
let j_scanning: i64 = self.read_key("jPointsAreConsecutive")?;
231+
232+
if ![0, 1].contains(&j_scanning) {
233+
return Err(MessageNdarrayError::UnexpectedKeyValue(
234+
"jPointsAreConsecutive".to_owned(),
235+
)
236+
.into());
237+
}
238+
239+
let j_scanning = j_scanning != 0;
240+
241+
let shape = if j_scanning {
242+
(ni, nj, 3_usize)
243+
} else {
244+
(nj, ni, 3_usize)
245+
};
246+
247+
let mut latlonvals =
248+
Array3::from_shape_vec(shape, latlonvals).map_err(MessageNdarrayError::from)?;
249+
250+
if j_scanning {
251+
latlonvals.swap_axes(0, 1);
252+
}
253+
254+
let (lats, lons, vals) =
255+
latlonvals
256+
.view_mut()
257+
.multi_slice_move((s![.., .., 0], s![.., .., 1], s![.., .., 2]));
258+
259+
Ok(RustyCodesMessage {
260+
longitudes: lons.into_owned(),
261+
latitudes: lats.into_owned(),
262+
values: vals.into_owned(),
263+
})
264+
}
265+
}
266+
141267
#[cfg(test)]
142268
mod tests {
143269
use fallible_iterator::FallibleIterator;
144270
use float_cmp::assert_approx_eq;
145271

146272
use super::*;
147-
use crate::DynamicKeyType;
148273
use crate::ProductKind;
149274
use crate::codes_handle::CodesHandle;
275+
use crate::keyed_message::DynamicKeyType;
150276
use std::path::Path;
151277

152278
#[test]

tests/handle.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
use std::{path::Path, thread};
22

33
use anyhow::{Context, Result};
4-
use eccodes::{CodesHandle, DynamicKeyType, FallibleIterator, ProductKind};
4+
use eccodes::{CodesHandle, keyed_message::DynamicKeyType, FallibleIterator, ProductKind};
55

66
#[test]
77
fn thread_safety() {

0 commit comments

Comments
 (0)