use crate::{
as_rc, client::QueryOptions, AsKeys, DataSignal, Fetcher, QueryClient, QueryData, Status,
};
use fluvio_wasm_timer::Delay;
use std::any::Any;
use std::{future::Future, rc::Rc};
use sycamore::{
futures::spawn_local,
reactive::{
create_effect, create_memo, create_rc_signal, create_ref, create_selector, use_context,
ReadSignal, Scope, Signal,
},
};
pub struct Query<'a, T, E, F: Fn()> {
pub data: &'a ReadSignal<QueryData<Rc<T>, Rc<E>>>,
pub status: Rc<Signal<Status>>,
pub refetch: &'a F,
}
impl QueryClient {
pub(crate) fn find_query(
&self,
key: &[u64],
new_hook: bool,
) -> Option<(Rc<DataSignal>, Rc<Signal<Status>>, Fetcher)> {
let data = self.data_signals.read().unwrap().get(key);
let status = self.status_signals.read().unwrap().get(key);
let fetcher = self.fetchers.read().unwrap().get(key)?.clone();
let (data, status) = match (data, status) {
(None, None) => None,
(None, Some(status)) => {
let data = if let Some(data) = self.cache.read().unwrap().get(key) {
QueryData::Ok(data)
} else {
QueryData::Loading
};
let data = as_rc(create_rc_signal(data));
if new_hook {
self.data_signals
.write()
.unwrap()
.insert(key.to_vec(), data.clone());
}
Some((data, status))
}
(Some(data), None) => {
let status = as_rc(create_rc_signal(Status::Success));
if new_hook {
self.status_signals
.write()
.unwrap()
.insert(key.to_vec(), status.clone());
}
Some((data, status))
}
(Some(data), Some(status)) => Some((data, status)),
}?;
Some((data, status, fetcher))
}
pub(crate) fn insert_query(
&self,
key: Vec<u64>,
data: Rc<DataSignal>,
status: Rc<Signal<Status>>,
fetcher: Fetcher,
) {
self.data_signals.write().unwrap().insert(key.clone(), data);
self.status_signals
.write()
.unwrap()
.insert(key.clone(), status);
self.fetchers.write().unwrap().insert(key, fetcher);
}
pub(crate) fn run_query(
self: Rc<Self>,
key: &[u64],
data: Rc<DataSignal>,
status: Rc<Signal<Status>>,
fetcher: Fetcher,
options: &QueryOptions,
) {
let options = self.default_options.merge(options);
if let Some(cached) = {
let cache = self.cache.read().unwrap();
cache.get(key)
} {
data.set(QueryData::Ok(cached));
self.clone().invalidate_queries(vec![key.to_vec()]);
} else if *status.get_untracked() != Status::Fetching {
status.set(Status::Fetching);
let key = key.to_vec();
spawn_local(async move {
let mut res = fetcher().await;
let mut retries = 0;
while res.is_err() && retries < options.retries {
Delay::new((options.retry_fn)(retries)).await.unwrap();
res = fetcher().await;
retries += 1;
}
data.set(res.map_or_else(QueryData::Err, QueryData::Ok));
if let QueryData::Ok(data) = data.get_untracked().as_ref() {
self.cache
.write()
.unwrap()
.insert(key, data.clone(), &options);
}
status.set(Status::Success);
});
}
}
pub(crate) fn refetch_query(self: Rc<Self>, key: &[u64]) {
self.invalidate_queries(vec![key.to_vec()]);
}
}
pub fn use_query<'a, K, T, E, F, R>(
cx: Scope<'a>,
key: K,
fetcher: F,
) -> Query<'a, T, E, impl Fn() + 'a>
where
K: AsKeys + 'a,
F: Fn() -> R + 'static,
R: Future<Output = Result<T, E>> + 'static,
T: 'static,
E: 'static,
{
use_query_with_options(cx, key, fetcher, QueryOptions::default())
}
pub fn use_query_with_options<'a, K, T, E, F, R>(
cx: Scope<'a>,
key: K,
fetcher: F,
options: QueryOptions,
) -> Query<'a, T, E, impl Fn() + 'a>
where
K: AsKeys + 'a,
F: Fn() -> R + 'static,
R: Future<Output = Result<T, E>> + 'static,
T: 'static,
E: 'static,
{
let id = create_selector(cx, move || key.as_keys());
let client = use_context::<Rc<QueryClient>>(cx).clone();
let (data, status, fetcher) = if let Some(query) = client.find_query(&id.get(), true) {
query
} else {
let data: Rc<DataSignal> = as_rc(create_rc_signal(QueryData::Loading));
let status = as_rc(create_rc_signal(Status::Idle));
let fetcher: Fetcher = Rc::new(move || {
let fut = fetcher();
Box::pin(async move {
fut.await
.map(|data| -> Rc<dyn Any> { Rc::new(data) })
.map_err(|err| -> Rc<dyn Any> { Rc::new(err) })
})
});
client.insert_query(
id.get().as_ref().clone(),
data.clone(),
status.clone(),
fetcher.clone(),
);
(data, status, fetcher)
};
{
let client = client.clone();
let data = data.clone();
let status = status.clone();
create_effect(cx, move || {
log::info!("Key changed. New key: {:?}", id.get());
client.clone().run_query(
&id.get(),
data.clone(),
status.clone(),
fetcher.clone(),
&options,
);
});
}
let refetch = create_ref(cx, move || {
client.clone().refetch_query(&id.get());
});
let data = create_memo(cx, move || match data.get().as_ref() {
QueryData::Loading => QueryData::Loading,
QueryData::Ok(data) => QueryData::Ok(data.clone().downcast().unwrap()),
QueryData::Err(err) => QueryData::Err(err.clone().downcast().unwrap()),
});
Query {
data,
status,
refetch,
}
}