use std::time::Duration;
use futures::future::select_all;
use libhaystack::val::kind::HaystackKind;
use super::sleep::sleep_millis;
use crate::base::Status;
use crate::base::block::{Block, BlockState, convert_value_kind};
use crate::base::input::{InputProps, input_reader::InputReader};
use crate::blocks::InputImpl;
use crate::blocks::utils::get_sleep_dur;
pub type ReaderImpl = <InputImpl as InputProps>::Reader;
pub type WriterImpl = <InputImpl as InputProps>::Writer;
impl<B: Block> InputReader for B {
async fn read_inputs(&mut self) -> Option<usize> {
read_block_inputs(self).await
}
async fn read_inputs_until_ready(&mut self) -> Option<usize> {
const MAX_BACKOFF_MS: u64 = 2000;
let mut backoff = get_sleep_dur();
loop {
let result = read_block_inputs(self).await;
if result.is_some() {
return result;
}
sleep_millis(backoff).await;
backoff = (backoff.saturating_mul(2)).min(MAX_BACKOFF_MS);
}
}
async fn wait_on_inputs(&mut self, timeout: Duration) -> Option<usize> {
let millis = timeout.as_millis() as u64;
tokio::select! {
_ = sleep_millis(millis) => None,
result = async {
loop {
match self.read_inputs().await {
Some(idx) => return Some(idx),
None => std::future::pending::<()>().await,
}
}
} => result,
}
}
}
pub(crate) async fn read_block_inputs<B: Block>(block: &mut B) -> Option<usize> {
if let Some(idx) = drain_ready_inputs(block) {
return Some(idx);
}
{
let mut connected: Vec<_> = block
.inputs_mut()
.into_iter()
.filter(|i| i.is_connected())
.collect();
if connected.is_empty() {
return None;
}
let futures: Vec<_> = connected.iter_mut().map(|i| i.receiver()).collect();
let _ = select_all(futures).await;
}
drain_ready_inputs(block)
}
fn drain_ready_inputs<B: Block>(block: &mut B) -> Option<usize> {
let mut last_idx = None;
let mut conversion_fault: Option<String> = None;
let mut upstream_fault: Option<String> = None;
{
let mut inputs = block.inputs_mut();
for (idx, input) in inputs.iter_mut().enumerate() {
if !input.is_connected() {
continue;
}
let Some((value, status)) = input.try_take() else {
continue;
};
let expected = *input.kind();
let actual = HaystackKind::from(&value);
let converted = if expected != HaystackKind::Null && expected != actual {
match convert_value_kind(value, expected, actual) {
Ok(v) => v,
Err(err) => {
log::error!("Error converting value: {}", err);
if conversion_fault.is_none() {
conversion_fault = Some(format!(
"type conversion failed on input {}: {}",
input.name(),
err
));
}
continue;
}
}
} else {
value
};
input.set_value(converted, status);
if status == Status::Fault && upstream_fault.is_none() {
upstream_fault = Some(format!("upstream fault on input {}", input.name()));
}
last_idx = Some(idx);
}
}
if let Some(reason) = conversion_fault {
block.set_state(BlockState::fault(reason));
} else if let Some(reason) = upstream_fault {
block.set_state(BlockState::fault(reason));
}
last_idx
}