apple-quant-algorithmic 0.3.0

Apple Quant's algorithmic library.
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;
	}
}