use std::{range::Range, str::FromStr};
use databento::dbn::{
encode::{
AsyncDbnEncoder, AsyncEncodeRecord, AsyncEncodeRecordRef, DynEncoderBuilder,
},
Dataset, MetadataBuilder, RecordRef, SType, Schema,
};
use time::UtcDateTime;
use tokio::{fs::File, io::BufWriter};
use crate::{
instrument::{InstrumentSpec, InstrumentTicker},
aggregation::TradeTradeTimestamp,
};
use super::symbology::Symbology;
pub struct DatabentoMarketDataEncoder {
dbn_encoder: AsyncDbnEncoder<BufWriter<File>>,
}
impl DatabentoMarketDataEncoder {
pub(super) async fn new(
first_buf_writer: BufWriter<File>,
date_time_range: &Range<UtcDateTime>,
instrument_ticker: &InstrumentTicker,
) -> Self {
let (
symbol,
stype,
) = instrument_ticker.symbol();
let metadata = MetadataBuilder::new()
.dataset(Dataset::GlbxMdp3)
.schema(Some(Schema::Trades))
.start(date_time_range.start.unix_timestamp_nanos() as u64)
.end(Some(
std::num::NonZeroU64::new(
date_time_range.end.unix_timestamp_nanos() as u64,
).unwrap(),
))
.symbols(vec![String::from_str(&symbol).unwrap(),])
.stype_in(Some(stype))
.stype_out(SType::InstrumentId)
.ts_out(true).build();
let mut dbn_encoder = AsyncDbnEncoder::new(
first_buf_writer,
&metadata,
).await.unwrap();
Self { dbn_encoder }
}
pub(super) async fn encode_record_ref(
&mut self,
record_ref: RecordRef<'_>,
) {
self.dbn_encoder.encode_record_ref(record_ref).await.unwrap();
}
pub(super) async fn shutdown(
mut self,
) {
self.dbn_encoder.flush().await;
self.dbn_encoder.shutdown().await;
}
}