use std::any::{Any, TypeId};
use std::collections::HashMap;
use std::sync::Arc;
use crate::{EffectsHandle, Memo, Step};
use taquba::LeaseHandle;
use tokio_util::sync::CancellationToken;
use crate::Result;
use crate::jobs::handle::JobHandle;
use crate::jobs::job::Job;
use crate::jobs::runner::{Inner, SubmitOptions};
#[derive(Default)]
pub(crate) struct State {
map: HashMap<TypeId, Arc<dyn Any + Send + Sync>>,
}
impl State {
pub(crate) fn insert<T: Any + Send + Sync>(&mut self, value: T) {
self.map.insert(TypeId::of::<T>(), Arc::new(value));
}
pub(crate) fn get<T: Any + Send + Sync>(&self) -> Option<&T> {
self.map
.get(&TypeId::of::<T>())
.and_then(|value| value.downcast_ref::<T>())
}
}
pub struct JobContext<'a> {
inner: Arc<Inner>,
step: &'a Step,
}
impl<'a> JobContext<'a> {
pub(crate) fn new(inner: Arc<Inner>, step: &'a Step) -> Self {
Self { inner, step }
}
pub fn state<T: Any + Send + Sync>(&self) -> &T {
self.try_state().unwrap_or_else(|| {
panic!(
"no application state of type `{}` registered on the JobRunner",
std::any::type_name::<T>()
)
})
}
pub fn try_state<T: Any + Send + Sync>(&self) -> Option<&T> {
self.inner.state().get::<T>()
}
pub fn id(&self) -> &str {
&self.step.run_id
}
pub fn attempt(&self) -> u32 {
self.step.attempts
}
pub fn cancel_token(&self) -> &CancellationToken {
&self.step.cancel_token
}
pub fn lease(&self) -> &LeaseHandle {
&self.step.lease
}
pub fn memo(&self) -> &Memo {
&self.step.memo
}
pub fn effects(&self) -> &EffectsHandle {
&self.step.effects
}
pub async fn kv_get(&self, key: &[u8]) -> Result<Option<Vec<u8>>> {
Ok(self.step.kv.get(key).await?.map(|bytes| bytes.to_vec()))
}
pub async fn submit<J: Job>(&self, job: J) -> Result<JobHandle<J>> {
self.inner.submit(job, SubmitOptions::default()).await
}
}