use std::future::{Future, IntoFuture};
use std::marker::PhantomData;
use std::pin::Pin;
use std::sync::Arc;
use std::time::Duration;
use crate::{RunStatus, StepErrorKind};
use thiserror::Error;
use crate::jobs::job::Job;
use crate::jobs::runner::Inner;
use crate::outcome::{
StoredErrorKind, StoredOutcome, Terminal, Unrecorded, read_outcome, wait_terminal,
};
use crate::{Error, Result};
const JOIN_CHUNK: Duration = Duration::from_secs(3600);
#[derive(Debug, Clone, Error)]
#[error("job failed ({kind:?}): {message}")]
pub struct JobError {
pub kind: StepErrorKind,
pub message: String,
}
#[derive(Debug, Error)]
pub enum JoinError {
#[error(transparent)]
Infra(#[from] Error),
#[error(transparent)]
Job(#[from] JobError),
}
pub struct JobHandle<J: Job> {
id: String,
queue_job_id: Option<String>,
inner: Arc<Inner>,
newly_submitted: bool,
_marker: PhantomData<fn() -> J>,
}
impl<J: Job> Clone for JobHandle<J> {
fn clone(&self) -> Self {
Self {
id: self.id.clone(),
queue_job_id: self.queue_job_id.clone(),
inner: self.inner.clone(),
newly_submitted: self.newly_submitted,
_marker: PhantomData,
}
}
}
impl<J: Job> JobHandle<J> {
pub(crate) fn new(
id: String,
queue_job_id: Option<String>,
inner: Arc<Inner>,
newly_submitted: bool,
) -> Self {
Self {
id,
queue_job_id,
inner,
newly_submitted,
_marker: PhantomData,
}
}
pub fn id(&self) -> &str {
&self.id
}
pub fn newly_submitted(&self) -> bool {
self.newly_submitted
}
pub async fn status(&self) -> Option<RunStatus> {
self.inner.runtime().status(&self.id).await
}
pub async fn fetch_result(&self) -> Result<Option<std::result::Result<J::Output, JobError>>> {
match read_outcome(&self.inner.run_memo(&self.id)).await? {
None => Ok(None),
Some(record) => Ok(Some(decode_outcome::<J>(record.outcome)?)),
}
}
pub async fn join(&self) -> Result<std::result::Result<J::Output, JobError>> {
loop {
if let Some(outcome) = self.join_timeout(JOIN_CHUNK).await? {
return Ok(outcome);
}
}
}
pub async fn join_timeout(
&self,
timeout: Duration,
) -> Result<Option<std::result::Result<J::Output, JobError>>> {
let Some(queue_job_id) = &self.queue_job_id else {
return match self.fetch_result().await? {
Some(outcome) => Ok(Some(outcome)),
None => Err(Error::JobNotFound(self.id.clone())),
};
};
let terminal = wait_terminal(
self.inner.queue(),
&self.inner.run_memo(&self.id),
queue_job_id,
timeout,
)
.await?;
match terminal {
None => Ok(None),
Some(Terminal::Recorded(record)) => Ok(Some(decode_outcome::<J>(record.outcome)?)),
Some(Terminal::Unrecorded(Unrecorded::NotFound)) => {
Err(Error::JobNotFound(self.id.clone()))
}
Some(Terminal::Unrecorded(unrecorded)) => {
let message = match unrecorded {
Unrecorded::Dead(Some(error)) => error,
_ => "job terminated without recording an outcome".to_string(),
};
Ok(Some(Err(JobError {
kind: StepErrorKind::Transient,
message,
})))
}
}
}
}
fn decode_outcome<J: Job>(
outcome: StoredOutcome,
) -> Result<std::result::Result<J::Output, JobError>> {
match outcome {
StoredOutcome::Success { output } => Ok(Ok(rmp_serde::from_slice(&output)?)),
StoredOutcome::Failure { kind, message } => Ok(Err(JobError {
kind: match kind {
StoredErrorKind::Transient => StepErrorKind::Transient,
StoredErrorKind::Permanent => StepErrorKind::Permanent,
},
message,
})),
}
}
impl<J: Job> IntoFuture for JobHandle<J> {
type Output = std::result::Result<J::Output, JoinError>;
type IntoFuture = Pin<Box<dyn Future<Output = Self::Output> + Send>>;
fn into_future(self) -> Self::IntoFuture {
Box::pin(async move {
match self.join().await {
Ok(Ok(output)) => Ok(output),
Ok(Err(job_error)) => Err(JoinError::Job(job_error)),
Err(infra) => Err(JoinError::Infra(infra)),
}
})
}
}