taskvisor 0.3.0

Event-driven task orchestration with restart, backoff, and user-defined subscribers
Documentation
//! Integration tests for `add_and_watch` / `TaskWaiter`.

mod common;

use std::time::Duration;

use common::*;
use taskvisor::prelude::*;

const ADD_TIMEOUT: Duration = Duration::from_secs(1);

fn supervisor() -> (std::sync::Arc<Supervisor>, SupervisorHandle) {
    let sup = Supervisor::new(SupervisorConfig::default(), vec![]);
    let handle = sup.serve();
    (sup, handle)
}

#[tokio::test]
async fn outcome_reason_is_byte_identical_to_the_event_reason() {
    use std::sync::Arc;

    let collector = EventCollector::new();
    let subs: Vec<Arc<dyn Subscribe>> = vec![collector.clone() as Arc<dyn Subscribe>];
    let sup = Supervisor::new(SupervisorConfig::default(), subs);
    let handle = sup.serve();

    let spec = TaskSpec::restartable(make_fail("drifter", Some(9)))
        .with_backoff(fast_backoff())
        .with_max_retries(2);
    let (id, waiter) = handle
        .add_and_watch(spec, ADD_TIMEOUT)
        .await
        .expect("add_and_watch should succeed");

    let outcome = with_timeout(5, waiter.wait())
        .await
        .expect("waiter errored");

    assert!(
        poll_until(Duration::from_secs(2), || async {
            collector
                .by_id(id)
                .iter()
                .any(|e| e.kind == EventKind::ActorExhausted)
        })
        .await
    );
    let event = collector
        .by_id(id)
        .into_iter()
        .find(|e| e.kind == EventKind::ActorExhausted)
        .expect("ActorExhausted event for the run");

    match outcome {
        TaskOutcome::Failed { reason, exit_code } => {
            assert_eq!(
                &*reason,
                event.reason.as_deref().expect("event carries a reason"),
                "TaskOutcome reason must be byte-identical to the ActorExhausted reason"
            );
            assert_eq!(exit_code, event.exit_code, "exit_code must match too");
        }
        other => panic!("expected Failed, got {other:?}"),
    }

    let _ = handle.shutdown().await;
}

#[tokio::test]
async fn completed_outcome_for_successful_once_task() {
    let (_sup, handle) = supervisor();

    let (id, waiter) = handle
        .add_and_watch(TaskSpec::once(make_ok_once("ok")), ADD_TIMEOUT)
        .await
        .expect("add_and_watch should succeed");
    assert_eq!(waiter.id(), id);

    let outcome = with_timeout(5, waiter.wait())
        .await
        .expect("waiter errored");
    assert_eq!(outcome, TaskOutcome::Completed);
    assert!(outcome.is_success());

    let _ = handle.shutdown().await;
}

#[tokio::test]
async fn failed_outcome_carries_reason_and_exit_code() {
    let (_sup, handle) = supervisor();

    let spec = TaskSpec::restartable(make_fail("flaky", Some(7)))
        .with_backoff(fast_backoff())
        .with_max_retries(2);
    let (_id, waiter) = handle
        .add_and_watch(spec, ADD_TIMEOUT)
        .await
        .expect("add_and_watch should succeed");

    match with_timeout(5, waiter.wait())
        .await
        .expect("waiter errored")
    {
        TaskOutcome::Failed { reason, exit_code } => {
            assert!(
                reason.contains("max_retries_exceeded"),
                "reason must mention exhausted retries: {reason}"
            );
            assert_eq!(exit_code, Some(7), "exit code must survive to the outcome");
        }
        other => panic!("expected Failed, got {other:?}"),
    }

    let _ = handle.shutdown().await;
}

#[tokio::test]
async fn fatal_outcome_for_fatal_error() {
    let (_sup, handle) = supervisor();

    let (_id, waiter) = handle
        .add_and_watch(
            TaskSpec::restartable(make_fatal("doomed", Some(137))),
            ADD_TIMEOUT,
        )
        .await
        .expect("add_and_watch should succeed");

    match with_timeout(5, waiter.wait())
        .await
        .expect("waiter errored")
    {
        TaskOutcome::Fatal { reason, exit_code } => {
            assert!(
                reason.contains("unrecoverable"),
                "reason must carry the fatal message: {reason}"
            );
            assert_eq!(exit_code, Some(137));
        }
        other => panic!("expected Fatal, got {other:?}"),
    }

    let _ = handle.shutdown().await;
}

#[tokio::test]
async fn failed_outcome_after_task_panic_with_never_policy() {
    let (_sup, handle) = supervisor();

    let (_id, waiter) = handle
        .add_and_watch(TaskSpec::once(make_panic("kaboom")), ADD_TIMEOUT)
        .await
        .expect("add_and_watch should succeed");

    match with_timeout(5, waiter.wait())
        .await
        .expect("waiter errored")
    {
        TaskOutcome::Failed { reason, .. } => {
            assert!(
                reason.contains("panic"),
                "reason must mention the panic: {reason}"
            );
        }
        other => panic!("expected Failed, got {other:?}"),
    }

    let _ = handle.shutdown().await;
}

#[tokio::test]
async fn spurious_canceled_return_resolves_canceled_outcome() {
    let (_sup, handle) = supervisor();

    let liar: TaskRef = TaskFn::arc("liar-watch", |_ctx: CancellationToken| async {
        Err(TaskError::Canceled)
    });
    let (_id, waiter) = handle
        .add_and_watch(TaskSpec::restartable(liar), ADD_TIMEOUT)
        .await
        .expect("add_and_watch should succeed");

    let outcome = with_timeout(5, waiter.wait())
        .await
        .expect("waiter errored");
    assert_eq!(
        outcome,
        TaskOutcome::Canceled,
        "a task returning Canceled without cancellation must resolve as Canceled"
    );

    let _ = handle.shutdown().await;
}

#[tokio::test]
async fn shutdown_drain_force_aborts_stubborn_watched_task() {
    let cfg = SupervisorConfig {
        grace: Duration::from_millis(150),
        ..Default::default()
    };
    let sup = Supervisor::new(cfg, vec![]);
    let handle = sup.serve();

    let stubborn: TaskRef = TaskFn::arc("stubborn-watch", |_ctx: CancellationToken| async {
        tokio::time::sleep(Duration::from_secs(60)).await;
        Ok(())
    });
    let (_id, waiter) = handle
        .add_and_watch(TaskSpec::once(stubborn), ADD_TIMEOUT)
        .await
        .expect("add_and_watch should succeed");

    let (shutdown_res, outcome) = tokio::join!(handle.shutdown(), with_timeout(5, waiter.wait()));
    assert!(
        shutdown_res.is_err(),
        "stubborn task must trip GraceExceeded"
    );
    assert_eq!(
        outcome.expect("waiter errored"),
        TaskOutcome::ForceAborted,
        "the shutdown drain's force-abort must resolve the waiter as ForceAborted"
    );
}

#[tokio::test]
async fn waiter_stays_pending_across_periodic_reruns() {
    let (_sup, handle) = supervisor();

    let spec =
        TaskSpec::restartable(make_ok_once("periodic-watch")).with_restart(RestartPolicy::Always {
            interval: Some(Duration::from_millis(20)),
        });
    let (id, waiter) = handle
        .add_and_watch(spec, ADD_TIMEOUT)
        .await
        .expect("add_and_watch should succeed");

    let pending = tokio::time::timeout(Duration::from_millis(200), waiter.wait()).await;
    assert!(
        pending.is_err(),
        "waiter must stay pending across successful Always re-runs"
    );

    let _ = handle.cancel(id).await;
    let _ = handle.shutdown().await;
}

#[tokio::test]
async fn cancelled_outcome_when_task_is_cancelled() {
    let (_sup, handle) = supervisor();

    let (id, waiter) = handle
        .add_and_watch(TaskSpec::restartable(make_coop("coop")), ADD_TIMEOUT)
        .await
        .expect("add_and_watch should succeed");

    let removed = handle.cancel(id).await.expect("cancel should not error");
    assert!(removed, "existing task must report removed=true");

    let outcome = with_timeout(5, waiter.wait())
        .await
        .expect("waiter errored");
    assert_eq!(outcome, TaskOutcome::Canceled);

    let _ = handle.shutdown().await;
}

#[tokio::test]
async fn force_aborted_outcome_for_noncooperative_task() {
    let cfg = SupervisorConfig {
        grace: Duration::from_millis(100),
        ..Default::default()
    };
    let sup = Supervisor::new(cfg, vec![]);
    let handle = sup.serve();

    let stubborn: TaskRef = TaskFn::arc("stubborn", |_ctx: CancellationToken| async move {
        tokio::time::sleep(Duration::from_secs(60)).await;
        Ok(())
    });
    let (id, waiter) = handle
        .add_and_watch(TaskSpec::once(stubborn), ADD_TIMEOUT)
        .await
        .expect("add_and_watch should succeed");

    handle.remove(id).expect("remove should be accepted");

    let outcome = with_timeout(5, waiter.wait())
        .await
        .expect("waiter errored");
    assert_eq!(outcome, TaskOutcome::ForceAborted);

    let _ = handle.shutdown().await;
}

#[tokio::test]
async fn duplicate_name_returns_already_exists_not_a_waiter() {
    let (_sup, handle) = supervisor();

    let first = handle
        .add_and_watch(TaskSpec::restartable(make_coop("dup")), ADD_TIMEOUT)
        .await;
    assert!(first.is_ok(), "first add must succeed");

    let second = handle
        .add_and_watch(TaskSpec::restartable(make_coop("dup")), ADD_TIMEOUT)
        .await;
    assert!(
        matches!(second, Err(RuntimeError::TaskAlreadyExists { .. })),
        "duplicate add must surface TaskAlreadyExists, got {second:?}"
    );

    let _ = handle.shutdown().await;
}

#[tokio::test]
async fn shutdown_resolves_pending_waiters() {
    let (_sup, handle) = supervisor();

    let (_id, waiter) = handle
        .add_and_watch(TaskSpec::restartable(make_coop("worker")), ADD_TIMEOUT)
        .await
        .expect("add_and_watch should succeed");

    handle
        .clone()
        .shutdown()
        .await
        .expect("shutdown should be Ok");

    let outcome = with_timeout(5, waiter.wait())
        .await
        .expect("waiter errored");
    assert_eq!(
        outcome,
        TaskOutcome::Canceled,
        "cooperative task must resolve as Canceled on shutdown"
    );
}

#[tokio::test]
async fn dropping_waiter_does_not_affect_task() {
    let (_sup, handle) = supervisor();

    let (id, waiter) = handle
        .add_and_watch(TaskSpec::restartable(make_coop("ignored")), ADD_TIMEOUT)
        .await
        .expect("add_and_watch should succeed");
    drop(waiter);

    assert!(
        poll_until(Duration::from_secs(2), || async {
            handle.is_alive("ignored").await
        })
        .await,
        "task must keep running after its waiter is dropped"
    );

    let removed = handle.cancel(id).await.expect("cancel should not error");
    assert!(removed);

    let _ = handle.shutdown().await;
}

#[tokio::test]
async fn outcome_is_delivered_even_under_bus_lag() {
    let cfg = SupervisorConfig {
        bus_capacity: 2,
        ..Default::default()
    };
    let sup = Supervisor::new(cfg, vec![]);
    let handle = sup.serve();

    let spec = TaskSpec::restartable(make_fail("noisy", None))
        .with_backoff(fast_backoff())
        .with_max_retries(5);
    let (_id, waiter) = handle
        .add_and_watch(spec, ADD_TIMEOUT)
        .await
        .expect("add_and_watch should succeed");

    match with_timeout(5, waiter.wait())
        .await
        .expect("waiter errored")
    {
        TaskOutcome::Failed { .. } => {}
        other => panic!("expected Failed despite bus lag, got {other:?}"),
    }

    let _ = handle.shutdown().await;
}