use std::collections::hash_map::Entry;
use std::collections::{HashMap, HashSet};
use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU32, AtomicU8, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
use arrow_array::RecordBatch;
use parking_lot::{Condvar, Mutex, RwLock};
use slab::Slab;
use tokio::sync::{mpsc, oneshot};
use xbbg_core::{
AsyncSession, AuthConfig, BlpError, CorrelationId, EntitlementCheck, EventType, SeatType,
};
const SERVICE_OPEN_TIMEOUT_MS: u64 = 10_000;
const IDENTITY_AUTH_TIMEOUT_MS: u64 = 10_000;
const APIAUTH_SERVICE: &str = "//blp/apiauth";
const SLOW_REQUEST_WARN_THRESHOLD: Duration = Duration::from_secs(30);
use super::dispatch::{DispatchKey, IDENTITY_CID, IDENTITY_TOKEN_CID, SERVICE_OPEN_CID_TAG};
use super::state::{
BqlState, BsrchState, BulkDataState, FieldInfoState, GenericState, HistDataState,
HistDataStreamState, IntradayBarState, IntradayBarStreamState, IntradayTickState,
IntradayTickStreamState, RefDataState,
};
use super::{
attach_auth_context, build_session_options, EngineConfig, PlannedRequestShape, PreparedRequest,
RequestParams, SlabKey, WorkerHealth, SESSION_STARTUP_TIMEOUT_MS,
};
fn iter_named_request_parameters(
params: &RequestParams,
) -> impl Iterator<Item = (&str, &str)> + '_ {
params
.elements
.iter()
.flat_map(|pairs| pairs.iter())
.chain(params.options.iter().flat_map(|pairs| pairs.iter()))
.map(|(name, value)| (name.as_str(), value.as_str()))
}
fn apply_named_request_parameter(
request: &mut xbbg_core::Request,
name: &str,
value: &str,
) -> Result<(), BlpError> {
if name.contains('.') {
if let Ok(int_val) = value.parse::<i32>() {
request.set_nested_int(name, int_val)?;
} else {
request.set_nested_str(name, value)?;
}
} else if request.set_str(name, value).is_err() {
request.append_str(name, value)?;
}
Ok(())
}
fn apply_excel_grid_request_parameters(
request: &mut xbbg_core::Request,
params: &RequestParams,
) -> Result<(), BlpError> {
if let Some(domain) = params
.elements
.iter()
.flat_map(|pairs| pairs.iter())
.rev()
.find(|(name, _)| name.eq_ignore_ascii_case("Domain"))
.map(|(_, value)| value.as_str())
{
request.set_str("Domain", domain)?;
}
let Some(overrides) = params
.overrides
.as_ref()
.filter(|values| !values.is_empty())
else {
return Ok(());
};
let overrides_ptr = request.get_or_create_element("Overrides")?;
for (name, value) in overrides {
if name.is_empty() {
continue;
}
let entry_ptr = unsafe { request.append_element(overrides_ptr)? };
unsafe { request.set_element_string(entry_ptr, "name", name)? };
unsafe { request.set_element_string(entry_ptr, "value", value)? };
}
Ok(())
}
#[allow(clippy::large_enum_variant)]
pub enum UnifiedRequestState {
RefData(RefDataState),
HistData(HistDataState),
BulkData(BulkDataState),
HistDataStream(HistDataStreamState),
Generic(GenericState),
Bql(BqlState),
Bsrch(BsrchState),
FieldInfo(FieldInfoState),
IntradayBar(IntradayBarState),
IntradayTick(IntradayTickState),
IntradayBarStream(IntradayBarStreamState),
IntradayTickStream(IntradayTickStreamState),
}
impl UnifiedRequestState {
pub fn on_partial(&mut self, msg: &xbbg_core::Message) {
match self {
UnifiedRequestState::RefData(s) => s.on_partial(msg),
UnifiedRequestState::HistData(s) => s.on_partial(msg),
UnifiedRequestState::BulkData(s) => s.on_partial(msg),
UnifiedRequestState::HistDataStream(s) => s.on_partial(msg),
UnifiedRequestState::Generic(s) => s.on_partial(msg),
UnifiedRequestState::Bql(s) => s.on_partial(msg),
UnifiedRequestState::Bsrch(s) => s.on_partial(msg),
UnifiedRequestState::FieldInfo(s) => s.on_partial(msg),
UnifiedRequestState::IntradayBar(s) => s.on_partial(msg),
UnifiedRequestState::IntradayTick(s) => s.on_partial(msg),
UnifiedRequestState::IntradayBarStream(s) => s.on_partial(msg),
UnifiedRequestState::IntradayTickStream(s) => s.on_partial(msg),
}
}
pub fn finish_and_reply(self, msg: &xbbg_core::Message) {
match self {
UnifiedRequestState::RefData(s) => s.finish(msg),
UnifiedRequestState::HistData(s) => s.finish(msg),
UnifiedRequestState::BulkData(s) => s.finish(msg),
UnifiedRequestState::HistDataStream(s) => s.finish(msg),
UnifiedRequestState::Generic(s) => s.finish(msg),
UnifiedRequestState::Bql(s) => s.finish(msg),
UnifiedRequestState::Bsrch(s) => s.finish(msg),
UnifiedRequestState::FieldInfo(s) => s.finish(msg),
UnifiedRequestState::IntradayBar(s) => s.finish(msg),
UnifiedRequestState::IntradayTick(s) => s.finish(msg),
UnifiedRequestState::IntradayBarStream(s) => s.finish(msg),
UnifiedRequestState::IntradayTickStream(s) => s.finish(msg),
}
}
pub fn fail(self, error: BlpError) {
match self {
UnifiedRequestState::RefData(s) => {
let _ = s.reply.send(Err(error));
}
UnifiedRequestState::HistData(s) => {
let _ = s.reply.send(Err(error));
}
UnifiedRequestState::BulkData(s) => {
let _ = s.reply.send(Err(error));
}
UnifiedRequestState::HistDataStream(s) => s.fail(error),
UnifiedRequestState::Generic(s) => {
let _ = s.reply.send(Err(error));
}
UnifiedRequestState::Bql(s) => {
let _ = s.reply.send(Err(error));
}
UnifiedRequestState::Bsrch(s) => {
let _ = s.reply.send(Err(error));
}
UnifiedRequestState::FieldInfo(s) => {
let _ = s.reply.send(Err(error));
}
UnifiedRequestState::IntradayBar(s) => {
let _ = s.reply.send(Err(error));
}
UnifiedRequestState::IntradayTick(s) => {
let _ = s.reply.send(Err(error));
}
UnifiedRequestState::IntradayBarStream(s) => s.fail(error),
UnifiedRequestState::IntradayTickStream(s) => s.fail(error),
}
}
}
fn send_stream_error(stream: mpsc::Sender<Result<RecordBatch, BlpError>>, error: BlpError) {
let _ = stream.blocking_send(Err(error));
}
#[derive(Clone, Copy, Debug)]
pub(crate) struct RequestTicket {
pub(crate) key: SlabKey,
pub(crate) generation: u32,
}
struct RequestStateCell {
generation: u32,
state: Mutex<Option<UnifiedRequestState>>,
active_partials: Mutex<usize>,
active_partials_cv: Condvar,
}
impl RequestStateCell {
fn new(generation: u32, state: UnifiedRequestState) -> Self {
Self {
generation,
state: Mutex::new(Some(state)),
active_partials: Mutex::new(0),
active_partials_cv: Condvar::new(),
}
}
fn begin_partial(self: &Arc<Self>) -> RequestPartialGuard {
*self.active_partials.lock() += 1;
RequestPartialGuard {
cell: Arc::clone(self),
}
}
fn wait_for_partials(&self) {
let mut active = self.active_partials.lock();
while *active != 0 {
self.active_partials_cv.wait(&mut active);
}
}
fn on_partial(&self, msg: &xbbg_core::Message<'_>) -> bool {
let mut state = self.state.lock();
let Some(state) = state.as_mut() else {
return false;
};
state.on_partial(msg);
true
}
fn finish_and_reply(&self, msg: &xbbg_core::Message<'_>) -> bool {
self.wait_for_partials();
let Some(state) = self.state.lock().take() else {
return false;
};
state.finish_and_reply(msg);
true
}
fn fail(&self, error: BlpError) -> bool {
self.wait_for_partials();
let Some(state) = self.state.lock().take() else {
return false;
};
state.fail(error);
true
}
fn discard(&self) -> bool {
self.wait_for_partials();
self.state.lock().take().is_some()
}
}
struct RequestPartialGuard {
cell: Arc<RequestStateCell>,
}
impl RequestPartialGuard {
fn on_partial(&self, msg: &xbbg_core::Message<'_>) -> bool {
self.cell.on_partial(msg)
}
}
impl Drop for RequestPartialGuard {
fn drop(&mut self) {
let mut active = self.cell.active_partials.lock();
debug_assert!(*active > 0);
*active = active.saturating_sub(1);
if *active == 0 {
self.cell.active_partials_cv.notify_all();
}
}
}
struct RequestSlot {
sent_at: Instant,
warned: bool,
cell: Arc<RequestStateCell>,
}
struct PendingServiceOpen {
cid: i64,
waiters: Vec<oneshot::Sender<Result<(), BlpError>>>,
}
#[derive(Default)]
enum IdentityAuthState {
#[default]
NotRequested,
Pending(Vec<oneshot::Sender<Result<(), BlpError>>>),
Ready,
Failed(String),
}
#[derive(Default)]
struct StartupLatch {
resolved: bool,
result: Option<Result<(), BlpError>>,
}
pub(super) struct WorkerShared {
id: usize,
requests: Mutex<Slab<RequestSlot>>,
next_generation: AtomicU32,
open_services: RwLock<HashSet<String>>,
pending_service_opens: Mutex<HashMap<String, PendingServiceOpen>>,
next_service_open_id: AtomicI64,
startup: Mutex<StartupLatch>,
startup_cv: Condvar,
identity: Mutex<IdentityAuthState>,
pending_token: Mutex<Option<std::sync::mpsc::SyncSender<Result<String, BlpError>>>>,
pending_identity_auth: Mutex<Option<std::sync::mpsc::SyncSender<Result<(), BlpError>>>>,
health: Arc<AtomicU8>,
shutting_down: AtomicBool,
}
impl WorkerShared {
fn new(id: usize, health: Arc<AtomicU8>) -> Self {
Self {
id,
requests: Mutex::new(Slab::new()),
next_generation: AtomicU32::new(0),
open_services: RwLock::new(HashSet::new()),
pending_service_opens: Mutex::new(HashMap::new()),
next_service_open_id: AtomicI64::new(0),
startup: Mutex::new(StartupLatch::default()),
startup_cv: Condvar::new(),
identity: Mutex::new(IdentityAuthState::default()),
pending_token: Mutex::new(None),
pending_identity_auth: Mutex::new(None),
health,
shutting_down: AtomicBool::new(false),
}
}
fn next_generation(&self) -> u32 {
self.next_generation
.fetch_add(1, Ordering::Relaxed)
.wrapping_add(1)
}
fn insert_slot(&self, generation: u32, state: UnifiedRequestState) -> SlabKey {
self.requests.lock().insert(RequestSlot {
sent_at: Instant::now(),
warned: false,
cell: Arc::new(RequestStateCell::new(generation, state)),
})
}
fn take_slot(&self, dispatch_key: DispatchKey) -> Option<RequestSlot> {
let key = dispatch_key.to_slab_key();
let mut requests = self.requests.lock();
match requests.get(key) {
Some(slot) if slot.cell.generation == dispatch_key.generation() => {
Some(requests.remove(key))
}
Some(_) => {
xbbg_log::debug!(
worker_id = self.id,
key = key,
"stale message for recycled slot; dropped"
);
None
}
None => None,
}
}
fn resolve_startup(&self, result: Result<(), BlpError>) {
let mut latch = self.startup.lock();
if !latch.resolved {
latch.resolved = true;
latch.result = Some(result);
self.startup_cv.notify_all();
}
}
fn wait_startup(&self, timeout: Duration) -> Result<(), BlpError> {
let deadline = Instant::now() + timeout;
let mut latch = self.startup.lock();
while latch.result.is_none() {
if self.startup_cv.wait_until(&mut latch, deadline).timed_out() {
return Err(BlpError::Timeout);
}
}
latch.result.take().expect("checked above")
}
pub(super) fn scan_timeouts(&self, hard_timeout: Option<Duration>) -> Vec<RequestTicket> {
let now = Instant::now();
let mut expired = Vec::new();
let mut requests = self.requests.lock();
for (key, slot) in requests.iter_mut() {
let elapsed = now.duration_since(slot.sent_at);
if elapsed > SLOW_REQUEST_WARN_THRESHOLD && !slot.warned {
slot.warned = true;
xbbg_log::warn!(
worker_id = self.id,
request_key = key,
elapsed_secs = elapsed.as_secs(),
"request outstanding longer than expected; still waiting on Bloomberg"
);
}
if let Some(hard) = hard_timeout {
if elapsed >= hard {
expired.push(RequestTicket {
key,
generation: slot.cell.generation,
});
}
}
}
expired
}
fn drain_in_flight(&self, reason: &str) {
let drained: Vec<RequestSlot> = {
let mut requests = self.requests.lock();
requests.drain().collect()
};
if drained.is_empty() {
return;
}
let count = drained.len();
for slot in drained {
slot.cell.fail(BlpError::Internal {
detail: format!("{} (worker={})", reason, self.id),
});
}
xbbg_log::error!(
worker_id = self.id,
drained = count,
reason = reason,
"drained in-flight requests"
);
}
fn fail_pending_service_opens(&self, reason: &str) {
let drained: Vec<(String, PendingServiceOpen)> = {
let mut pending = self.pending_service_opens.lock();
pending.drain().collect()
};
for (service, open) in drained {
for waiter in open.waiters {
let _ = waiter.send(Err(BlpError::OpenService {
service: service.clone(),
source: None,
label: Some(reason.to_string()),
}));
}
}
}
fn resolve_identity(&self, outcome: Result<(), BlpError>) {
let detail = outcome.as_ref().err().map(ToString::to_string);
let waiters = {
let mut state = self.identity.lock();
let next = match &detail {
None => IdentityAuthState::Ready,
Some(e) => IdentityAuthState::Failed(e.clone()),
};
match std::mem::replace(&mut *state, next) {
IdentityAuthState::Pending(waiters) => waiters,
_ => Vec::new(),
}
};
for waiter in waiters {
let _ = waiter.send(match &detail {
None => Ok(()),
Some(e) => Err(BlpError::Internal { detail: e.clone() }),
});
}
}
fn fail_pending_identity(&self, reason: &str) {
let is_pending = matches!(&*self.identity.lock(), IdentityAuthState::Pending(_));
if is_pending {
self.resolve_identity(Err(BlpError::Internal {
detail: format!("identity authorization aborted: {reason}"),
}));
}
if let Some(tx) = self.pending_token.lock().take() {
let _ = tx.send(Err(BlpError::Internal {
detail: format!("identity token generation aborted: {reason}"),
}));
}
if let Some(tx) = self.pending_identity_auth.lock().take() {
let _ = tx.send(Err(BlpError::Internal {
detail: format!("identity authorization aborted: {reason}"),
}));
}
}
fn resolve_identity_outcome(&self, outcome: Result<(), BlpError>) {
let detail = outcome.as_ref().err().map(ToString::to_string);
if let Some(tx) = self.pending_identity_auth.lock().take() {
let _ = tx.send(match &detail {
None => Ok(()),
Some(e) => Err(BlpError::Internal { detail: e.clone() }),
});
}
self.resolve_identity(outcome);
}
fn handle_token_status(&self, msg: &xbbg_core::Message<'_>) {
let ours = msg.correlation_ids().any(|cid| {
matches!(
cid,
CorrelationId::Int(value) if value == IDENTITY_TOKEN_CID
)
});
if !ours {
return;
}
match msg.type_str() {
"TokenGenerationSuccess" => {
let token = msg
.elements()
.get_by_str("token")
.and_then(|e| e.get_str(0))
.map(str::to_string);
if let Some(tx) = self.pending_token.lock().take() {
let _ = tx.send(token.ok_or_else(|| BlpError::Internal {
detail: "TokenGenerationSuccess carried no token element".to_string(),
}));
}
}
"TokenGenerationFailure" => {
let reason = extract_reason_description(msg)
.unwrap_or_else(|| "TokenGenerationFailure".to_string());
xbbg_log::warn!(
worker_id = self.id,
reason = reason.as_str(),
"identity token generation failed"
);
if let Some(tx) = self.pending_token.lock().take() {
let _ = tx.send(Err(BlpError::Internal {
detail: format!("identity token generation failed: {reason}"),
}));
}
}
other => {
xbbg_log::debug!(
worker_id = self.id,
message_type = other,
"unhandled token status message"
);
}
}
}
fn handle_identity_authorization_message(&self, msg: &xbbg_core::Message<'_>) {
match msg.type_str() {
"AuthorizationSuccess" => {
xbbg_log::info!(worker_id = self.id, "identity authorized");
self.resolve_identity_outcome(Ok(()));
}
"AuthorizationFailure" => {
let reason = extract_reason_description(msg);
xbbg_log::warn!(
worker_id = self.id,
reason = %reason.as_deref().unwrap_or(""),
"identity authorization failed"
);
self.resolve_identity_outcome(Err(BlpError::Internal {
detail: format!(
"identity authorization failed: {}",
reason.unwrap_or_else(|| "AuthorizationFailure".to_string())
),
}));
}
"AuthorizationRevoked" => {
let reason = extract_reason_description(msg);
xbbg_log::warn!(
worker_id = self.id,
reason = %reason.as_deref().unwrap_or(""),
"identity authorization revoked"
);
*self.identity.lock() = IdentityAuthState::Failed(format!(
"identity authorization revoked: {}",
reason.unwrap_or_default()
));
}
other => {
xbbg_log::debug!(
worker_id = self.id,
message_type = other,
"unhandled authorization message"
);
}
}
}
fn is_identity_auth_message(msg: &xbbg_core::Message<'_>) -> bool {
msg.correlation_ids().any(|cid| {
matches!(
cid,
CorrelationId::Int(value) if value == IDENTITY_CID
)
})
}
fn handle_authorization_status(&self, msg: &xbbg_core::Message<'_>) {
if Self::is_identity_auth_message(msg) {
self.handle_identity_authorization_message(msg);
}
}
pub(super) fn dispatch_event(&self, ev: xbbg_core::Event) {
let et = ev.event_type();
for msg in ev.iter() {
match et {
EventType::PartialResponse => {
self.handle_partial_response(&msg);
}
EventType::Response => {
self.handle_response(&msg);
}
EventType::RequestStatus => {
self.handle_request_status(&msg);
}
EventType::SessionStatus => {
self.handle_session_status(&msg);
}
EventType::ServiceStatus => {
self.handle_service_status(&msg);
}
EventType::AuthorizationStatus => {
self.handle_authorization_status(&msg);
}
EventType::TokenStatus => {
self.handle_token_status(&msg);
}
_ => {}
}
}
}
fn handle_partial_response(&self, msg: &xbbg_core::Message<'_>) {
if Self::is_identity_auth_message(msg) {
self.handle_identity_authorization_message(msg);
return;
}
for correlation_id in msg.correlation_ids() {
let Some(dispatch_key) = DispatchKey::from_correlation_id(&correlation_id) else {
continue;
};
let key = dispatch_key.to_slab_key();
let partial = {
let requests = self.requests.lock();
match requests.get(key) {
Some(slot) if slot.cell.generation == dispatch_key.generation() => {
Some(slot.cell.begin_partial())
}
Some(_) => {
xbbg_log::debug!(
worker_id = self.id,
key = key,
"stale partial response for recycled slot; dropped"
);
None
}
None => None,
}
};
if let Some(partial) = partial {
if partial.on_partial(msg) {
xbbg_log::trace!(worker_id = self.id, key = key, "partial response");
}
}
}
}
fn handle_response(&self, msg: &xbbg_core::Message<'_>) {
if Self::is_identity_auth_message(msg) {
self.handle_identity_authorization_message(msg);
return;
}
for correlation_id in msg.correlation_ids() {
let Some(dispatch_key) = DispatchKey::from_correlation_id(&correlation_id) else {
continue;
};
if let Some(slot) = self.take_slot(dispatch_key) {
let key = dispatch_key.to_slab_key();
let rtt_ms = slot.sent_at.elapsed().as_micros() as f64 / 1000.0;
xbbg_log::debug!(
worker_id = self.id,
rtt_ms = rtt_ms,
key = key,
"bloomberg_roundtrip"
);
if slot.cell.finish_and_reply(msg) {
xbbg_log::debug!(worker_id = self.id, key = key, "response completed");
}
}
}
}
fn handle_request_status(&self, msg: &xbbg_core::Message<'_>) {
let msg_type = msg.type_str();
if msg_type != "RequestFailure" {
return;
}
if Self::is_identity_auth_message(msg) {
let reason =
extract_reason_description(msg).unwrap_or_else(|| "RequestFailure".to_string());
self.resolve_identity_outcome(Err(BlpError::Internal {
detail: format!("identity authorization request failed: {reason}"),
}));
return;
}
for correlation_id in msg.correlation_ids() {
let Some(dispatch_key) = DispatchKey::from_correlation_id(&correlation_id) else {
continue;
};
let reason = extract_reason_description(msg);
xbbg_log::error!(
worker_id = self.id,
key = dispatch_key.to_slab_key(),
reason = %reason.as_deref().unwrap_or(""),
"request failed"
);
if let Some(slot) = self.take_slot(dispatch_key) {
slot.cell.fail(BlpError::Internal {
detail: reason.unwrap_or_else(|| "RequestFailure".to_string()),
});
}
}
}
fn handle_session_status(&self, msg: &xbbg_core::Message<'_>) {
let msg_type = msg.type_str();
match msg_type {
"SessionStarted" => {
self.health.store(0, Ordering::Release);
self.resolve_startup(Ok(()));
xbbg_log::info!(worker_id = self.id, "session started");
}
"SessionStartupFailure" => {
let reason = extract_reason_description(msg);
xbbg_log::error!(
worker_id = self.id,
reason = %reason.as_deref().unwrap_or(""),
"session startup failure"
);
self.health.store(2, Ordering::Release);
self.resolve_startup(Err(session_start_error("session startup failure", reason)));
}
"SessionTerminated" => {
let reason = extract_reason_description(msg);
self.resolve_startup(Err(session_start_error(
"session terminated during startup",
reason.clone(),
)));
self.drain_in_flight(reason.as_deref().unwrap_or("Bloomberg session terminated"));
self.fail_pending_service_opens("Bloomberg session terminated");
self.fail_pending_identity("Bloomberg session terminated");
self.health.store(2, Ordering::Release);
if self.shutting_down.load(Ordering::Acquire) {
xbbg_log::info!(
worker_id = self.id,
reason = %reason.as_deref().unwrap_or(""),
"SessionTerminated during requested shutdown"
);
} else {
xbbg_log::error!(
worker_id = self.id,
reason = %reason.as_deref().unwrap_or(""),
"SessionTerminated — worker is dead"
);
}
}
"AuthorizationFailure" => {
let reason = extract_reason_description(msg);
self.resolve_startup(Err(session_start_error(
"session identity authorization failed",
reason,
)));
}
"AuthorizationRevoked" => {
let reason = extract_reason_description(msg);
self.resolve_startup(Err(session_start_error(
"session identity authorization revoked",
reason.clone(),
)));
self.drain_in_flight(
reason
.as_deref()
.unwrap_or("Bloomberg session identity revoked"),
);
self.fail_pending_service_opens("Bloomberg session identity revoked");
self.fail_pending_identity("Bloomberg session identity revoked");
self.health.store(2, Ordering::Release);
xbbg_log::error!(
worker_id = self.id,
reason = %reason.as_deref().unwrap_or(""),
"AuthorizationRevoked — identity gone; worker is dead"
);
}
"SessionConnectionDown" => {
let reason = extract_reason_description(msg);
self.drain_in_flight(
reason
.as_deref()
.unwrap_or("Bloomberg session connection lost (transient)"),
);
self.health.store(1, Ordering::Release);
if self.shutting_down.load(Ordering::Acquire) {
xbbg_log::info!(
worker_id = self.id,
reason = %reason.as_deref().unwrap_or(""),
"SessionConnectionDown during requested shutdown"
);
} else {
xbbg_log::warn!(
worker_id = self.id,
reason = %reason.as_deref().unwrap_or(""),
"SessionConnectionDown — failing in-flight requests; SDK will auto-reconnect"
);
}
}
"SessionConnectionUp" => {
self.health.store(0, Ordering::Release);
let reason = extract_reason_description(msg);
xbbg_log::info!(
worker_id = self.id,
reason = %reason.as_deref().unwrap_or(""),
"SessionConnectionUp — worker healthy"
);
}
_ => {
xbbg_log::debug!(worker_id = self.id, msg_type = msg_type, "session status");
}
}
}
fn handle_service_status(&self, msg: &xbbg_core::Message<'_>) {
let msg_type = msg.type_str();
if matches!(msg_type, "ServiceOpened" | "ServiceOpenFailure") {
if let Some(CorrelationId::Int(cid_int)) = msg.correlation_id(0) {
let entry = {
let mut pending = self.pending_service_opens.lock();
let service = pending
.iter()
.find_map(|(name, open)| (open.cid == cid_int).then(|| name.clone()));
service.and_then(|name| pending.remove(&name).map(|open| (name, open)))
};
if let Some((service, open)) = entry {
if msg_type == "ServiceOpened" {
self.open_services.write().insert(service.clone());
xbbg_log::debug!(worker_id = self.id, service = %service, "service opened");
for waiter in open.waiters {
let _ = waiter.send(Ok(()));
}
} else {
let reason = extract_reason_description(msg);
xbbg_log::warn!(
worker_id = self.id,
service = %service,
reason = %reason.as_deref().unwrap_or(""),
"service open failed"
);
for waiter in open.waiters {
let _ = waiter.send(Err(BlpError::OpenService {
service: service.clone(),
source: None,
label: reason.clone(),
}));
}
}
return;
}
}
}
xbbg_log::debug!(worker_id = self.id, msg_type = msg_type, "service status");
}
}
fn session_start_error(context: &str, reason: Option<String>) -> BlpError {
BlpError::SessionStart {
source: None,
label: Some(match reason {
Some(reason) => format!("{context}: {reason}"),
None => context.to_string(),
}),
}
}
fn hint_entitlement_requirements(err: BlpError) -> BlpError {
match err {
BlpError::Internal { detail } => BlpError::Internal {
detail: format!(
"{detail} (identity/entitlement checks require EMRS enrollment for \
Desktop API terminals, or SAPI/B-PIPE authorization via EngineConfig.auth)"
),
},
other => other,
}
}
pub(crate) struct AsyncRequestWorker {
pub(crate) id: usize,
session: AsyncSession,
shared: Arc<WorkerShared>,
config: Arc<EngineConfig>,
identity_flow_lock: tokio::sync::Mutex<()>,
}
impl AsyncRequestWorker {
pub(crate) fn new(id: usize, config: Arc<EngineConfig>) -> Result<Self, BlpError> {
let options = build_session_options(&config, false)?;
let health = Arc::new(AtomicU8::new(0));
let shared = Arc::new(WorkerShared::new(id, health));
let handler_shared = Arc::clone(&shared);
let session = AsyncSession::new(&options, move |event| {
handler_shared.dispatch_event(event);
})?;
session
.start()
.map_err(|err| attach_auth_context(err, config.auth.as_ref()))?;
shared
.wait_startup(Duration::from_millis(u64::from(SESSION_STARTUP_TIMEOUT_MS)))
.map_err(|err| attach_auth_context(err, config.auth.as_ref()))?;
let worker = Self {
id,
session,
shared,
config,
identity_flow_lock: tokio::sync::Mutex::new(()),
};
worker.warmup();
xbbg_log::info!(worker_id = id, "AsyncRequestWorker started");
Ok(worker)
}
fn warmup(&self) {
for service_name in &self.config.warmup_services {
match self.session.open_service(service_name) {
Ok(()) => {
self.shared
.open_services
.write()
.insert(service_name.clone());
}
Err(e) => {
xbbg_log::warn!(
worker_id = self.id,
service = %service_name,
error = %e,
"failed to pre-warm service"
);
}
}
}
xbbg_log::info!(
worker_id = self.id,
services = ?self.shared.open_services.read().iter().collect::<Vec<_>>(),
"worker pre-warmed"
);
}
pub(crate) fn health(&self) -> WorkerHealth {
match self.shared.health.load(Ordering::Acquire) {
0 => WorkerHealth::Healthy,
1 => WorkerHealth::Degraded,
_ => WorkerHealth::Dead,
}
}
async fn ensure_service(&self, name: &str) -> Result<(), BlpError> {
if self.shared.open_services.read().contains(name) {
return Ok(());
}
let rx = {
let mut pending = self.shared.pending_service_opens.lock();
if self.shared.open_services.read().contains(name) {
return Ok(());
}
let (tx, rx) = oneshot::channel();
match pending.entry(name.to_string()) {
Entry::Occupied(mut entry) => entry.get_mut().waiters.push(tx),
Entry::Vacant(entry) => {
let id = self
.shared
.next_service_open_id
.fetch_add(1, Ordering::Relaxed)
.wrapping_add(1);
let cid_int = SERVICE_OPEN_CID_TAG | (id & (SERVICE_OPEN_CID_TAG - 1));
let cid = CorrelationId::Int(cid_int);
self.session.open_service_async(name, &cid)?;
entry.insert(PendingServiceOpen {
cid: cid_int,
waiters: vec![tx],
});
}
}
rx
};
match tokio::time::timeout(Duration::from_millis(SERVICE_OPEN_TIMEOUT_MS), rx).await {
Ok(Ok(result)) => result,
Ok(Err(_)) => Err(BlpError::Internal {
detail: format!("service open for {name} dropped without resolution"),
}),
Err(_) => Err(BlpError::Timeout),
}
}
async fn ensure_identity(&self, auth_config: &AuthConfig) -> Result<(), BlpError> {
let rx = {
let mut state = self.shared.identity.lock();
match &mut *state {
IdentityAuthState::Ready => return Ok(()),
IdentityAuthState::Pending(waiters) => {
let (tx, rx) = oneshot::channel();
waiters.push(tx);
rx
}
state @ (IdentityAuthState::NotRequested | IdentityAuthState::Failed(_)) => {
if let IdentityAuthState::Failed(prior) = state {
xbbg_log::debug!(
worker_id = self.id,
prior_failure = prior.as_str(),
"retrying identity authorization"
);
}
let auth_options = auth_config.build_auth_options()?;
self.session.generate_authorized_identity_async(
&auth_options,
&CorrelationId::Int(IDENTITY_CID),
)?;
let (tx, rx) = oneshot::channel();
*state = IdentityAuthState::Pending(vec![tx]);
xbbg_log::debug!(
worker_id = self.id,
auth_method = auth_config.method_name(),
"identity authorization started"
);
rx
}
}
};
match tokio::time::timeout(Duration::from_millis(IDENTITY_AUTH_TIMEOUT_MS), rx).await {
Ok(Ok(result)) => result,
Ok(Err(_)) => Err(BlpError::Internal {
detail: "identity authorization dropped without resolution".to_string(),
}),
Err(_) => {
self.shared
.fail_pending_identity("identity authorization timed out");
Err(BlpError::Timeout)
}
}
}
fn recv_classic_step<T>(
&self,
rx: &std::sync::mpsc::Receiver<Result<T, BlpError>>,
pending: &Mutex<Option<std::sync::mpsc::SyncSender<Result<T, BlpError>>>>,
step: &str,
) -> Result<T, BlpError> {
match rx.recv_timeout(Duration::from_millis(IDENTITY_AUTH_TIMEOUT_MS)) {
Ok(result) => result,
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
pending.lock().take();
Err(BlpError::Timeout)
}
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => Err(BlpError::Internal {
detail: format!("identity {step} dropped without resolution"),
}),
}
}
fn classic_authorize_blocking<T>(
&self,
op: impl FnOnce(&Self, &xbbg_core::Identity) -> Result<T, BlpError>,
) -> Result<T, BlpError> {
let token_rx = {
let mut pending = self.shared.pending_token.lock();
let (tx, rx) = std::sync::mpsc::sync_channel(1);
*pending = Some(tx);
if let Err(err) = self
.session
.generate_token(&CorrelationId::Int(IDENTITY_TOKEN_CID))
{
pending.take();
return Err(err);
}
rx
};
let token = self
.recv_classic_step(&token_rx, &self.shared.pending_token, "token generation")
.map_err(hint_entitlement_requirements)?;
let auth_service = self.session.get_service(APIAUTH_SERVICE)?;
let mut request = auth_service.create_request("AuthorizationRequest")?;
request.set_str("token", &token)?;
let mut identity = self.session.create_identity()?;
let auth_rx = {
let mut pending = self.shared.pending_identity_auth.lock();
let (tx, rx) = std::sync::mpsc::sync_channel(1);
*pending = Some(tx);
if let Err(err) = self.session.send_authorization_request(
&request,
&mut identity,
&CorrelationId::Int(IDENTITY_CID),
) {
pending.take();
return Err(err);
}
rx
};
self.recv_classic_step(
&auth_rx,
&self.shared.pending_identity_auth,
"authorization",
)
.map_err(hint_entitlement_requirements)?;
op(self, &identity)
}
async fn with_identity<T: Send + 'static>(
self: &Arc<Self>,
op: impl FnOnce(&Self, &xbbg_core::Identity) -> Result<T, BlpError> + Send + 'static,
) -> Result<T, BlpError> {
match &self.config.auth {
Some(auth_config) => {
self.ensure_identity(auth_config).await?;
let identity = self
.session
.authorized_identity(&CorrelationId::Int(IDENTITY_CID))?;
op(self, &identity)
}
None => {
self.ensure_service(APIAUTH_SERVICE).await?;
let _flow = self.identity_flow_lock.lock().await;
let this = Arc::clone(self);
tokio::task::spawn_blocking(move || this.classic_authorize_blocking(op))
.await
.map_err(|err| BlpError::Internal {
detail: format!("identity task failed to complete: {err}"),
})?
}
}
}
pub(crate) async fn identity_seat_type(self: &Arc<Self>) -> Result<SeatType, BlpError> {
self.with_identity(|_, identity| identity.seat_type()).await
}
pub(crate) async fn identity_check_entitlements(
self: &Arc<Self>,
service: &str,
eids: &[i32],
) -> Result<EntitlementCheck, BlpError> {
self.ensure_service(service).await?;
let service = service.to_string();
let eids = eids.to_vec();
self.with_identity(move |worker, identity| {
let service = worker.session.get_service(&service)?;
identity.check_entitlements(&service, &eids)
})
.await
}
pub(crate) async fn identity_is_authorized(
self: &Arc<Self>,
service: &str,
) -> Result<bool, BlpError> {
self.ensure_service(service).await?;
let service = service.to_string();
self.with_identity(move |worker, identity| {
let service = worker.session.get_service(&service)?;
Ok(identity.is_authorized(&service))
})
.await
}
pub(crate) async fn submit(
&self,
request: PreparedRequest,
reply: oneshot::Sender<Result<RecordBatch, BlpError>>,
) -> Option<RequestTicket> {
let params = request.params();
let t0 = Instant::now();
if let Err(error) = self.ensure_service(¶ms.service).await {
let _ = reply.send(Err(error));
return None;
}
xbbg_log::debug!(
worker_id = self.id,
elapsed_us = t0.elapsed().as_micros(),
"ensure_service"
);
xbbg_log::debug!(
worker_id = self.id,
shape = ?request.shape(),
fields = ?params.fields,
"creating request state"
);
let state = create_request_state(&request, reply);
self.dispatch_prepared(request, state)
}
pub(crate) async fn submit_stream(
&self,
request: PreparedRequest,
stream: mpsc::Sender<Result<RecordBatch, BlpError>>,
) -> Option<RequestTicket> {
let params = request.params();
let fields = params.fields.clone().unwrap_or_default();
let ticker = params.security.clone().unwrap_or_default();
let state = match request.shape() {
PlannedRequestShape::HistData(_) => {
UnifiedRequestState::HistDataStream(HistDataStreamState::new(fields, stream))
}
PlannedRequestShape::IntradayBar => {
UnifiedRequestState::IntradayBarStream(IntradayBarStreamState::new(ticker, stream))
}
PlannedRequestShape::IntradayTick => UnifiedRequestState::IntradayTickStream(
IntradayTickStreamState::new(ticker, stream),
),
_ => {
send_stream_error(
stream,
BlpError::InvalidArgument {
detail: format!(
"Streaming not supported for extractor: {:?}",
params.extractor
),
},
);
return None;
}
};
if let Err(error) = self.ensure_service(¶ms.service).await {
state.fail(error);
return None;
}
self.dispatch_prepared(request, state)
}
fn dispatch_prepared(
&self,
request: PreparedRequest,
state: UnifiedRequestState,
) -> Option<RequestTicket> {
let generation = self.shared.next_generation();
let key = self.shared.insert_slot(generation, state);
let dispatch_key = DispatchKey::with_generation(key, generation);
let cid = dispatch_key.to_correlation_id();
let result = (|| -> Result<(), BlpError> {
let params = request.params();
let service = self.session.get_service(¶ms.service)?;
xbbg_log::debug!(
worker_id = self.id,
operation = %request.effective_operation(),
securities = ?params.securities,
start_date = ?params.start_date,
end_date = ?params.end_date,
"building request"
);
let blp_request = build_request_from_params(self.id, &service, &request)?;
let actual_cid = self.session.send_request_with_label(
&blp_request,
Some(&cid),
params.request_id.as_deref(),
)?;
let actual_dispatch_key =
DispatchKey::from_correlation_id(&actual_cid).ok_or_else(|| {
BlpError::Internal {
detail: format!(
"Bloomberg returned non-dispatch correlation ID for request: {:?}",
actual_cid
),
}
})?;
if actual_dispatch_key != dispatch_key {
return Err(BlpError::Internal {
detail: format!(
"Bloomberg returned unexpected dispatch correlation ID {:?} for slab key {}",
actual_cid, key
),
});
}
xbbg_log::debug!(
worker_id = self.id,
key = key,
service = %params.service,
operation = %request.effective_operation(),
"request sent"
);
Ok(())
})();
match result {
Ok(()) => Some(RequestTicket { key, generation }),
Err(err) => {
if let Some(slot) = self.shared.take_slot(dispatch_key) {
slot.cell.fail(err);
}
None
}
}
}
pub(crate) fn cancel_request(&self, ticket: RequestTicket) {
let dispatch_key = DispatchKey::with_generation(ticket.key, ticket.generation);
let Some(slot) = self.shared.take_slot(dispatch_key) else {
return;
};
let cid = dispatch_key.to_correlation_id();
if let Err(error) = self.session.cancel(&cid) {
xbbg_log::warn!(
worker_id = self.id,
key = ticket.key,
error = %error,
"failed to cancel Bloomberg request"
);
slot.cell.discard();
self.shared.health.store(2, Ordering::Release);
self.shared
.drain_in_flight("Bloomberg request cancellation failed");
return;
}
slot.cell.discard();
xbbg_log::info!(
worker_id = self.id,
key = ticket.key,
"cancelled Bloomberg request"
);
}
pub(crate) fn timeout_request(&self, ticket: RequestTicket, timeout_ms: u64) {
let dispatch_key = DispatchKey::with_generation(ticket.key, ticket.generation);
let Some(slot) = self.shared.take_slot(dispatch_key) else {
return;
};
let cid = dispatch_key.to_correlation_id();
if let Err(error) = self.session.cancel(&cid) {
xbbg_log::warn!(
worker_id = self.id,
key = ticket.key,
error = %error,
"timeout_request: Bloomberg cancel failed"
);
}
slot.cell.fail(BlpError::Timeout);
xbbg_log::warn!(
worker_id = self.id,
key = ticket.key,
timeout_ms = timeout_ms,
"request exceeded request_timeout_ms; failed with BlpError::Timeout"
);
}
pub(crate) fn scan_timeouts(&self, hard_timeout: Option<Duration>) -> Vec<RequestTicket> {
self.shared.scan_timeouts(hard_timeout)
}
pub(crate) async fn introspect_schema(
&self,
service_uri: &str,
) -> Result<crate::schema::ServiceSchema, BlpError> {
xbbg_log::debug!(worker_id = self.id, service = %service_uri, "introspecting schema");
self.ensure_service(service_uri).await?;
let service = self.session.get_service(service_uri)?;
let schema = crate::schema::introspect_service(&service, service_uri);
xbbg_log::debug!(
worker_id = self.id,
service = %service_uri,
operations = schema.operations.len(),
"schema introspection complete"
);
Ok(schema)
}
pub(crate) fn signal_shutdown(&self) {
self.shared.shutting_down.store(true, Ordering::Release);
self.session.stop_async();
}
pub(crate) fn shutdown_blocking(&self) {
self.shared.shutting_down.store(true, Ordering::Release);
xbbg_log::info!(worker_id = self.id, "AsyncRequestWorker shutting down");
self.session.stop();
self.shared.fail_pending_service_opens("worker shutdown");
self.shared.fail_pending_identity("worker shutdown");
self.shared.drain_in_flight("worker shutdown");
}
}
fn create_request_state(
request: &PreparedRequest,
reply: oneshot::Sender<Result<RecordBatch, BlpError>>,
) -> UnifiedRequestState {
let params = request.params();
let fields = params.fields.clone().unwrap_or_default();
let field_types = params.field_types.clone();
match request.shape() {
PlannedRequestShape::RefData(output) => {
UnifiedRequestState::RefData(RefDataState::with_format(
fields,
output.format,
output.long_mode,
field_types,
params.include_security_errors,
reply,
))
}
PlannedRequestShape::HistData(output) => UnifiedRequestState::HistData(
HistDataState::with_format(fields, output.format, output.long_mode, field_types, reply),
),
PlannedRequestShape::BulkData => {
let field = fields.first().cloned().unwrap_or_default();
UnifiedRequestState::BulkData(BulkDataState::new(field, reply))
}
PlannedRequestShape::Generic => UnifiedRequestState::Generic(GenericState::new(reply)),
PlannedRequestShape::Bql => UnifiedRequestState::Bql(BqlState::new(reply)),
PlannedRequestShape::Bsrch => UnifiedRequestState::Bsrch(BsrchState::new(reply)),
PlannedRequestShape::FieldInfo => {
UnifiedRequestState::FieldInfo(FieldInfoState::new(reply))
}
PlannedRequestShape::IntradayBar => {
let ticker = params.security.clone().unwrap_or_default();
let event_type = params
.event_type
.clone()
.unwrap_or_else(|| "TRADE".to_string());
let interval = params.interval.unwrap_or(1);
UnifiedRequestState::IntradayBar(IntradayBarState::new(
ticker, event_type, interval, reply,
))
}
PlannedRequestShape::IntradayTick => {
let ticker = params.security.clone().unwrap_or_default();
UnifiedRequestState::IntradayTick(IntradayTickState::new(ticker, reply))
}
}
}
fn build_request_from_params(
worker_id: usize,
service: &xbbg_core::Service<'_>,
prepared: &PreparedRequest,
) -> Result<xbbg_core::Request, BlpError> {
let params = prepared.params();
let operation = prepared.effective_operation();
xbbg_log::trace!(operation = %operation, "creating request");
let mut request = service.create_request(operation)?;
xbbg_log::trace!("request created");
if prepared.is_excel_get_grid_request() {
apply_excel_grid_request_parameters(&mut request, params)?;
return Ok(request);
}
if let Some(securities) = ¶ms.securities {
for sec in securities {
xbbg_log::trace!(element = "securities", value = %sec, "appending");
request.append_str("securities", sec)?;
}
}
if let Some(security) = ¶ms.security {
if prepared.uses_intraday_security_element() {
xbbg_log::trace!(element = "security", value = %security, "setting scalar");
request.set_str("security", security)?;
} else {
xbbg_log::trace!(element = "securities", value = %security, "appending");
request.append_str("securities", security)?;
}
}
if let Some(fields) = ¶ms.fields {
for field in fields {
xbbg_log::trace!(element = "fields", value = %field, "appending");
request.append_str("fields", field)?;
}
}
if let Some(start) = ¶ms.start_date {
xbbg_log::trace!(element = "startDate", value = %start, "setting");
request.set_str("startDate", start)?;
}
if let Some(end) = ¶ms.end_date {
xbbg_log::trace!(element = "endDate", value = %end, "setting");
request.set_str("endDate", end)?;
}
if let Some(start) = ¶ms.start_datetime {
xbbg_log::trace!(element = "startDateTime", value = %start, "setting datetime");
request.set_datetime("startDateTime", start)?;
}
if let Some(end) = ¶ms.end_datetime {
xbbg_log::trace!(element = "endDateTime", value = %end, "setting datetime");
request.set_datetime("endDateTime", end)?;
}
if let Some(event_type) = ¶ms.event_type {
request.set_str("eventType", event_type)?;
}
if let Some(event_types) = ¶ms.event_types {
for et in event_types {
xbbg_log::trace!(element = "eventTypes", value = %et, "appending event type");
request.append_str("eventTypes", et)?;
}
}
if let Some(interval) = params.interval {
request.set_int("interval", interval as i32)?;
}
for (name, value) in iter_named_request_parameters(params) {
apply_named_request_parameter(&mut request, name, value)?;
}
if params.return_eids {
xbbg_log::trace!(element = "returnEids", "setting");
apply_named_request_parameter(&mut request, "returnEids", "true")?;
}
if let Some(field_ids) = ¶ms.field_ids {
for id in field_ids {
request.append_str("id", id)?;
}
}
if let Some(overrides) = ¶ms.overrides {
if !overrides.is_empty() {
let overrides_ptr = request.get_or_create_element("overrides")?;
for (field_id, value) in overrides {
let entry_ptr = unsafe { request.append_element(overrides_ptr)? };
unsafe { request.set_element_string(entry_ptr, "fieldId", field_id)? };
unsafe { request.set_element_string(entry_ptr, "value", value)? };
}
xbbg_log::debug!(
worker_id = worker_id,
count = overrides.len(),
"overrides applied"
);
}
}
if let Some(search_spec) = ¶ms.search_spec {
request.set_str("searchSpec", search_spec)?;
}
Ok(request)
}
fn extract_reason_description(msg: &xbbg_core::Message<'_>) -> Option<String> {
let reason = msg.elements().get_by_str("reason")?;
for key in ["description", "category", "message"] {
if let Some(s) = reason.get_by_str(key).and_then(|e| e.get_str(0)) {
return Some(s.to_string());
}
}
None
}
#[cfg(test)]
mod tests {
use super::*;
use crate::engine::RequestParams;
fn shared() -> WorkerShared {
WorkerShared::new(0, Arc::new(AtomicU8::new(0)))
}
fn generic_state() -> (
UnifiedRequestState,
oneshot::Receiver<Result<RecordBatch, BlpError>>,
) {
let (tx, rx) = oneshot::channel();
(UnifiedRequestState::Generic(GenericState::new(tx)), rx)
}
#[test]
fn take_slot_requires_matching_generation() {
let shared = shared();
let (state, _rx) = generic_state();
let generation = shared.next_generation();
let key = shared.insert_slot(generation, state);
let stale = DispatchKey::with_generation(key, generation.wrapping_add(1));
assert!(shared.take_slot(stale).is_none());
assert_eq!(shared.requests.lock().len(), 1);
let fresh = DispatchKey::with_generation(key, generation);
assert!(shared.take_slot(fresh).is_some());
assert!(shared.requests.lock().is_empty());
assert!(shared.take_slot(fresh).is_none());
}
#[test]
fn recycled_slot_is_not_visible_to_old_ticket() {
let shared = shared();
let (first, _rx1) = generic_state();
let g1 = shared.next_generation();
let key = shared.insert_slot(g1, first);
let first_key = DispatchKey::with_generation(key, g1);
assert!(shared.take_slot(first_key).is_some());
let (second, _rx2) = generic_state();
let g2 = shared.next_generation();
let key2 = shared.insert_slot(g2, second);
assert_eq!(key, key2, "slab should recycle the slot");
assert!(shared.take_slot(first_key).is_none());
assert_eq!(shared.requests.lock().len(), 1);
}
#[test]
fn scan_timeouts_warns_once_and_reports_expired() {
let shared = shared();
let (state, _rx) = generic_state();
let generation = shared.next_generation();
let key = shared.insert_slot(generation, state);
shared.requests.lock().get_mut(key).unwrap().sent_at =
Instant::now() - Duration::from_secs(120);
let expired = shared.scan_timeouts(Some(Duration::from_secs(60)));
assert_eq!(expired.len(), 1);
assert_eq!(expired[0].key, key);
assert_eq!(expired[0].generation, generation);
assert!(shared.requests.lock().get(key).unwrap().warned);
let expired = shared.scan_timeouts(None);
assert!(expired.is_empty());
}
#[test]
fn drain_in_flight_fails_all_slots() {
let shared = shared();
let (s1, mut rx1) = generic_state();
let (s2, mut rx2) = generic_state();
let g1 = shared.next_generation();
let g2 = shared.next_generation();
shared.insert_slot(g1, s1);
shared.insert_slot(g2, s2);
shared.drain_in_flight("test drain");
assert!(shared.requests.lock().is_empty());
assert!(matches!(rx1.try_recv(), Ok(Err(BlpError::Internal { .. }))));
assert!(matches!(rx2.try_recv(), Ok(Err(BlpError::Internal { .. }))));
}
#[test]
fn startup_latch_resolves_once() {
let shared = shared();
shared.resolve_startup(Ok(()));
shared.resolve_startup(Err(BlpError::Timeout)); assert!(shared.wait_startup(Duration::from_millis(10)).is_ok());
}
#[test]
fn startup_latch_times_out_when_unresolved() {
let shared = shared();
assert!(matches!(
shared.wait_startup(Duration::from_millis(10)),
Err(BlpError::Timeout)
));
}
#[test]
fn iter_named_request_parameters_includes_options_after_elements() {
let params = RequestParams {
elements: Some(vec![(
"periodicitySelection".to_string(),
"DAILY".to_string(),
)]),
options: Some(vec![("returnEids".to_string(), "true".to_string())]),
..Default::default()
};
let collected: Vec<(String, String)> = iter_named_request_parameters(¶ms)
.map(|(name, value)| (name.to_string(), value.to_string()))
.collect();
assert_eq!(
collected,
vec![
("periodicitySelection".to_string(), "DAILY".to_string()),
("returnEids".to_string(), "true".to_string()),
]
);
}
}