hans 0.1.0

Task orchestrator.
Documentation
use std::{future::Future, pin::Pin, sync::Arc};

use anyhow::Result;
use chrono::{DateTime, Utc};
use tokio::{spawn, sync::Mutex, task::JoinHandle};

use crate::{Executor, Task, TaskStatus};

/// A Tokio executor.
pub struct TokioExecutor;

impl Executor for TokioExecutor {}

/// An implementation of task for Tokio.
pub struct TokioTask {
    closure: Arc<Box<TokioTaskClosure>>,
    finish_status: Mutex<Option<TaskStatus>>,
    job: Mutex<Option<Job>>,
}

impl TokioTask {
    /// Creates new task.
    ///
    /// - `closure`: the closure to run.
    pub fn new<
        FN: Fn() -> FUT + Send + Sync + 'static,
        FUT: Future<Output = anyhow::Result<()>> + Send + Sync,
    >(
        closure: FN,
    ) -> Self {
        let closure = Arc::new(closure);
        Self {
            closure: Arc::new(Box::new(move || {
                let closure = closure.clone();
                Box::pin(async move { closure().await })
            })),
            finish_status: Mutex::new(None),
            job: Mutex::new(None),
        }
    }
}

impl Task<TokioExecutor> for TokioTask {
    async fn delete(&self, _exec: &TokioExecutor) -> Result<()> {
        Ok(())
    }

    async fn is_deleted(&self, _exec: &TokioExecutor) -> Result<bool> {
        Ok(true)
    }

    async fn start(&self, _exec: &TokioExecutor) -> Result<TaskStatus> {
        let closure = self.closure.clone();
        let mut job = self.job.lock().await;
        let started_at = Utc::now();
        *job = Some(Job {
            handle: spawn(async move {
                match closure().await {
                    Ok(()) => TaskStatus::Succeeded {
                        finished_at: Utc::now(),
                        started_at,
                    },
                    Err(_) => TaskStatus::Failed {
                        finished_at: Utc::now(),
                        started_at,
                    },
                }
            }),
            started_at,
        });
        Ok(TaskStatus::Running {
            started_at: Utc::now(),
        })
    }

    async fn status(&self, _exec: &TokioExecutor) -> Result<TaskStatus> {
        let mut finish_status = self.finish_status.lock().await;
        let mut guard = self.job.lock().await;
        if let Some(finish_status) = *finish_status {
            Ok(finish_status)
        } else if let Some(job) = guard.take() {
            if job.handle.is_finished() {
                let status = job.handle.await?;
                *finish_status = Some(status);
                Ok(status)
            } else {
                let started_at = job.started_at;
                *guard = Some(job);
                Ok(TaskStatus::Running { started_at })
            }
        } else {
            Ok(TaskStatus::Pending)
        }
    }
}

type TokioTaskClosure =
    dyn Fn() -> Pin<Box<dyn Future<Output = Result<()>> + Send + Sync>> + Send + Sync;

struct Job {
    handle: JoinHandle<TaskStatus>,
    started_at: DateTime<Utc>,
}