use async_trait::async_trait;
use crate::meta_storage::MetaBackend;
use crate::scheduler::events::PeerEventSenders;
use crate::scheduler::{SchedulerError, TaskRebuilder};
use crate::schema::{
CheckMode, JobLike, ResolutionLike, SharedProgress, TaskMetadata, TicketExplosion, TicketLike,
};
use crate::service::OperonService;
use crate::storage::OperonStorage;
#[derive(Debug)]
pub struct SpecWithMetadata<Svc, Sto, TS, const N: usize> {
pub spec: TS,
pub task_meta: TaskMetadata<N>,
_phantom: std::marker::PhantomData<(Svc, Sto)>,
}
impl<Svc, Sto, TS, const N: usize> SpecWithMetadata<Svc, Sto, TS, N> {
pub fn new(spec: TS, task_meta: TaskMetadata<N>) -> Self {
Self {
spec,
task_meta,
_phantom: std::marker::PhantomData,
}
}
}
impl<Svc, Sto, TS, const N: usize> Clone for SpecWithMetadata<Svc, Sto, TS, N>
where
TS: Clone,
{
fn clone(&self) -> Self {
Self {
spec: self.spec.clone(),
task_meta: self.task_meta,
_phantom: std::marker::PhantomData,
}
}
}
#[async_trait]
pub trait TaskSpec<Svc, Sto, MSto>: Clone + Send + Sync + 'static
where
Svc: OperonService,
Sto: OperonStorage,
MSto: MetaBackend,
{
type Job: JobLike;
type Resolution: ResolutionLike;
type Ticket: TicketLike;
type PeerEventSenders: PeerEventSenders<Svc::JobEnum, Svc::ResolutionEnum, Svc::TicketEnum>;
fn all_upstream_tasks(&self) -> Vec<&'static str>;
fn pool_size(&self) -> usize;
fn default_ticket(&self) -> Self::Ticket;
async fn check_consistency(
&self,
storage: &Sto,
client: MSto::Client<'_>,
mode: CheckMode,
) -> Result<bool, SchedulerError<Svc::Error, Sto::Error, MSto::Error>>;
async fn prepare_rebuild(
&self,
storage: &Sto,
progress: SharedProgress,
client: MSto::Client<'_>,
) -> Result<
Box<dyn TaskRebuilder<Svc, Sto, MSto>>,
SchedulerError<Svc::Error, Sto::Error, MSto::Error>,
>;
async fn run_job(
&self,
service: &Svc,
storage: &Sto,
meta: MSto,
job: Self::Job,
) -> Result<Self::Resolution, SchedulerError<Svc::Error, Sto::Error, MSto::Error>>;
async fn send_on_finish(
&self,
peer_txs: &Self::PeerEventSenders,
job: Self::Job,
resolution: Self::Resolution,
) -> Result<(), SchedulerError<Svc::Error, Sto::Error, MSto::Error>>;
async fn on_receive_job(
&self,
client: MSto::Client<'_>,
job: Svc::JobEnum,
) -> Result<Vec<Self::Ticket>, SchedulerError<Svc::Error, Sto::Error, MSto::Error>>;
async fn on_receive_resolution(
&self,
client: MSto::Client<'_>,
peer_txs: &Self::PeerEventSenders,
resolution: Svc::ResolutionEnum,
) -> Result<Vec<Self::Ticket>, SchedulerError<Svc::Error, Sto::Error, MSto::Error>>;
async fn on_receive_explosion(
&self,
client: MSto::Client<'_>,
explosion: TicketExplosion<Svc::TicketEnum>,
) -> Result<Vec<Self::Ticket>, SchedulerError<Svc::Error, Sto::Error, MSto::Error>>;
}