use std::sync::Arc;
use arrow_array::builder::StringBuilder;
use arrow_array::RecordBatch;
use arrow_schema::{DataType, Field, Schema};
use tokio::sync::oneshot;
use xbbg_log::trace;
use super::value_utils::top_level_response_error;
use crate::field_cache::BlpFieldType;
use xbbg_core::{BlpError, Message};
pub struct FieldInfoState {
field_builder: StringBuilder,
type_builder: StringBuilder,
description_builder: StringBuilder,
category_builder: StringBuilder,
pub reply: oneshot::Sender<Result<RecordBatch, BlpError>>,
}
impl FieldInfoState {
pub fn new(reply: oneshot::Sender<Result<RecordBatch, BlpError>>) -> Self {
Self {
field_builder: StringBuilder::new(),
type_builder: StringBuilder::new(),
description_builder: StringBuilder::new(),
category_builder: StringBuilder::new(),
reply,
}
}
pub fn on_partial(&mut self, msg: &Message) {
self.process_message(msg);
}
pub fn finish(mut self, msg: &Message) {
if let Some(error) = top_level_response_error(msg, "//blp/apiflds", "FieldInfoRequest") {
let _ = self.reply.send(Err(error));
return;
}
self.process_message(msg);
let result = self.build_batch();
if let Ok(ref batch) = result {
xbbg_log::debug!(rows = batch.num_rows(), "fieldinfo finish");
}
let _ = self.reply.send(result);
}
fn process_message(&mut self, msg: &Message) {
let root = msg.elements();
let Some(field_data) = root.get_by_str("fieldData") else {
trace!("No fieldData in message");
return;
};
let n = field_data.len();
for i in 0..n {
let Some(item) = field_data.get_element(i) else {
continue;
};
let Some(field_info) = item.get_by_str("fieldInfo") else {
continue;
};
let field = field_info
.get_by_str("mnemonic")
.and_then(|e| e.get_str(0))
.or_else(|| field_info.get_by_str("id").and_then(|e| e.get_str(0)))
.unwrap_or("");
if field.is_empty() {
continue;
}
let datatype = field_info.get_by_str("datatype");
let ftype = field_info.get_by_str("ftype");
let blp_type = BlpFieldType::from_metadata(
datatype.as_ref().and_then(|element| element.get_str(0)),
ftype.as_ref().and_then(|element| element.get_str(0)),
);
let arrow_type = blp_type.to_arrow_type_str();
let description = field_info
.get_by_str("description")
.and_then(|e| e.get_str(0))
.unwrap_or("");
let category = field_info
.get_by_str("categoryName")
.and_then(|cats| cats.get_str(0))
.unwrap_or("");
self.field_builder.append_value(field);
self.type_builder.append_value(arrow_type);
self.description_builder.append_value(description);
self.category_builder.append_value(category);
}
}
fn build_batch(&mut self) -> Result<RecordBatch, BlpError> {
let field_array = self.field_builder.finish();
let type_array = self.type_builder.finish();
let description_array = self.description_builder.finish();
let category_array = self.category_builder.finish();
let schema = Arc::new(Schema::new(vec![
Field::new("field", DataType::Utf8, false),
Field::new("type", DataType::Utf8, false),
Field::new("description", DataType::Utf8, true),
Field::new("category", DataType::Utf8, true),
]));
RecordBatch::try_new(
schema,
vec![
Arc::new(field_array),
Arc::new(type_array),
Arc::new(description_array),
Arc::new(category_array),
],
)
.map_err(|e| BlpError::Internal {
detail: format!("build RecordBatch: {e}"),
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::field_cache::FieldTypeResolver;
use arrow_array::StringArray;
use xbbg_core::test_support::TestEvent;
use xbbg_core::EventType;
#[test]
fn apiflds_datetime_ftype_reaches_the_field_cache() {
let schema = r#"<ServiceDefinition name="xbbg.test.fieldinfo" version="1.0.0.0">
<service name="//xbbg/test/fieldinfo" version="1.0.0.0">
<event name="FieldInfoResponse" eventType="FieldInfoResponseType"/>
</service>
<schema>
<sequenceType name="FieldInfoResponseType">
<element name="fieldData" type="FieldInfoDataType" minOccurs="0" maxOccurs="unbounded"/>
</sequenceType>
<sequenceType name="FieldInfoDataType">
<element name="fieldInfo" type="FieldMetadataType"/>
</sequenceType>
<sequenceType name="FieldMetadataType">
<element name="mnemonic" type="String"/>
<element name="datatype" type="String" minOccurs="0"/>
<element name="ftype" type="String" minOccurs="0"/>
</sequenceType>
</schema>
</ServiceDefinition>"#;
let event = TestEvent::with_schema(
schema,
EventType::Response,
"FieldInfoResponse",
&[],
|formatter| {
formatter.json(r#"{"fieldData":[
{"fieldInfo":{"mnemonic":"SYNTHETIC_TIME","datatype":"Datetime","ftype":"Time"}},
{"fieldInfo":{"mnemonic":"SYNTHETIC_DATE","datatype":"Datetime","ftype":"Date"}},
{"fieldInfo":{"mnemonic":"SYNTHETIC_STAMP","datatype":"Datetime","ftype":"Datetime"}},
{"fieldInfo":{"mnemonic":"SYNTHETIC_MIXED","datatype":"Datetime","ftype":"DateOrTime"}},
{"fieldInfo":{"mnemonic":"SYNTHETIC_PRICE","datatype":"Double","ftype":"Price"}},
{"fieldInfo":{"mnemonic":"SYNTHETIC_FALLBACK","ftype":"Time"}}
]}"#);
},
);
let (sender, mut receiver) = oneshot::channel();
FieldInfoState::new(sender).finish(&event.event().messages().next().unwrap());
let batch = receiver.try_recv().unwrap().unwrap();
let types = batch
.column_by_name("type")
.unwrap()
.as_any()
.downcast_ref::<StringArray>()
.unwrap();
assert_eq!(
types.iter().collect::<Vec<_>>(),
[
Some("time64"),
Some("date32"),
Some("timestamp"),
Some("string"),
Some("float64"),
Some("time64")
]
);
let dir = tempfile::tempdir().unwrap();
let resolver = FieldTypeResolver::with_cache_path(dir.path().join("field_cache.json"));
resolver.insert_from_response(&batch);
assert_eq!(
resolver.get_arrow_type("SYNTHETIC_TIME").as_deref(),
Some("time64")
);
assert_eq!(
resolver.get_arrow_type("SYNTHETIC_DATE").as_deref(),
Some("date32")
);
assert_eq!(
resolver.get_arrow_type("SYNTHETIC_STAMP").as_deref(),
Some("timestamp")
);
}
}