nautilus_persistence/python/wranglers/delta.rs
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92
// -------------------------------------------------------------------------------------------------
// Copyright (C) 2015-2024 Nautech Systems Pty Ltd. All rights reserved.
// https://nautechsystems.io
//
// Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
// You may not use this file except in compliance with the License.
// You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
// -------------------------------------------------------------------------------------------------
use std::{collections::HashMap, io::Cursor, str::FromStr};
use datafusion::arrow::ipc::reader::StreamReader;
use nautilus_core::python::to_pyvalue_err;
use nautilus_model::{data::delta::OrderBookDelta, identifiers::InstrumentId};
use pyo3::prelude::*;
use crate::arrow::DecodeFromRecordBatch;
#[pyclass()]
pub struct OrderBookDeltaDataWrangler {
instrument_id: InstrumentId,
price_precision: u8,
size_precision: u8,
metadata: HashMap<String, String>,
}
#[pymethods]
impl OrderBookDeltaDataWrangler {
#[new]
fn py_new(instrument_id: &str, price_precision: u8, size_precision: u8) -> PyResult<Self> {
let instrument_id = InstrumentId::from_str(instrument_id).map_err(to_pyvalue_err)?;
let metadata =
OrderBookDelta::get_metadata(&instrument_id, price_precision, size_precision);
Ok(Self {
instrument_id,
price_precision,
size_precision,
metadata,
})
}
#[getter]
fn instrument_id(&self) -> String {
self.instrument_id.to_string()
}
#[getter]
fn price_precision(&self) -> u8 {
self.price_precision
}
#[getter]
fn size_precision(&self) -> u8 {
self.size_precision
}
fn process_record_batch_bytes(
&self,
_py: Python,
data: &[u8],
) -> PyResult<Vec<OrderBookDelta>> {
// Create a StreamReader (from Arrow IPC)
let cursor = Cursor::new(data);
let reader = match StreamReader::try_new(cursor, None) {
Ok(reader) => reader,
Err(e) => return Err(to_pyvalue_err(e)),
};
let mut deltas = Vec::new();
// Read the record batches
for maybe_batch in reader {
let record_batch = match maybe_batch {
Ok(record_batch) => record_batch,
Err(e) => return Err(to_pyvalue_err(e)),
};
let batch_deltas = OrderBookDelta::decode_batch(&self.metadata, record_batch)
.map_err(to_pyvalue_err)?;
deltas.extend(batch_deltas);
}
Ok(deltas)
}
}