use std::error::Error;
use serde::de::DeserializeOwned;
use serde_json::Value;
use super::error::RuntimeError;
use super::phase::{Task, TaskDisplayKind, TaskId, TaskKey};
use super::reporting::{ActivityTask, TaskProgress};
use crate::configuration::TaskConfig;
pub type TaskResult = Result<(), Box<dyn Error + Send + Sync + 'static>>;
pub(crate) type Workload = Box<dyn FnOnce(&TaskContext) -> TaskResult + Send + 'static>;
enum TaskHandle {
Progress(TaskProgress),
Activity(ActivityTask),
}
pub struct TaskContext {
task: Task,
handle: TaskHandle,
}
impl TaskContext {
pub(crate) fn progress(task: Task, progress: TaskProgress) -> Self {
Self {
task,
handle: TaskHandle::Progress(progress),
}
}
pub(crate) fn activity(task: Task, activity: ActivityTask) -> Self {
Self {
task,
handle: TaskHandle::Activity(activity),
}
}
pub fn task(&self) -> &Task {
&self.task
}
pub fn key(&self) -> &TaskKey {
self.task.key()
}
pub fn id(&self) -> &TaskId {
self.task.id()
}
pub fn kind(&self) -> &str {
self.task.kind()
}
pub fn configuration(&self) -> &TaskConfig {
self.task.configuration()
}
pub fn value(&self, key: &str) -> Option<&Value> {
self.task.value(key)
}
pub fn decode_value<T>(&self, key: &str) -> Result<T, RuntimeError>
where
T: DeserializeOwned,
{
self.task.decode_value(key)
}
pub fn progress_handle(&self) -> Option<&TaskProgress> {
match &self.handle {
TaskHandle::Progress(progress) => Some(progress),
TaskHandle::Activity(_) => None,
}
}
pub fn set_target_iteration(&self, target: u64) -> Result<(), RuntimeError> {
self.required_progress()?.set_target_iteration(target)
}
pub fn set_iteration(&self, iteration: u64) -> Result<(), RuntimeError> {
self.required_progress()?.set_iteration(iteration)
}
pub fn should_continue(&self, iteration: u64) -> Result<bool, RuntimeError> {
self.required_progress()?.should_continue(iteration)
}
pub fn activity_handle(&self) -> Option<&ActivityTask> {
match &self.handle {
TaskHandle::Progress(_) => None,
TaskHandle::Activity(activity) => Some(activity),
}
}
pub fn is_cancelled(&self) -> bool {
match &self.handle {
TaskHandle::Progress(progress) => progress.is_cancelled(),
TaskHandle::Activity(activity) => activity.is_cancelled(),
}
}
pub fn set_detail(&self, detail: impl Into<String>) {
let detail = detail.into();
match &self.handle {
TaskHandle::Progress(progress) => progress.set_detail(detail),
TaskHandle::Activity(activity) => activity.set_detail(detail),
}
}
pub fn report(&self, message: impl Into<String>) -> Result<(), RuntimeError> {
let message = message.into();
match &self.handle {
TaskHandle::Progress(progress) => progress.report(message),
TaskHandle::Activity(activity) => activity.report(message),
}
}
pub(crate) fn complete(self) -> Result<(), RuntimeError> {
match self.handle {
TaskHandle::Progress(progress) => progress.complete(None),
TaskHandle::Activity(activity) => {
activity.complete();
Ok(())
}
}
}
pub(crate) fn fail(self, reason: impl Into<String>) {
let reason = reason.into();
match self.handle {
TaskHandle::Progress(progress) => progress.fail(reason),
TaskHandle::Activity(activity) => activity.fail(reason),
}
}
pub(crate) fn cancel(self, reason: impl Into<String>) {
let reason = reason.into();
match self.handle {
TaskHandle::Progress(progress) => progress.cancel(reason),
TaskHandle::Activity(activity) => activity.cancel(reason),
}
}
fn required_progress(&self) -> Result<&TaskProgress, RuntimeError> {
self.progress_handle()
.ok_or_else(|| RuntimeError::ManagedTaskKindMismatch {
task: self.key().to_string(),
requested: "progress",
actual: match self.task.display_kind() {
TaskDisplayKind::Progress => "progress",
TaskDisplayKind::Activity => "activity",
},
})
}
}
impl std::fmt::Debug for TaskContext {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter
.debug_struct("TaskContext")
.field("key", self.key())
.field("kind", &self.kind())
.field("display_kind", &self.task.display_kind())
.finish_non_exhaustive()
}
}