use std::sync::Arc;
use aion::registry::UnrecoverableRun;
use aion_core::WorkflowId;
use axum::body::Body;
use axum::http::{Request, StatusCode};
use chrono::{DateTime, TimeDelta, Utc};
use tower::ServiceExt;
use super::super::router::workflow_router;
use super::super::test_support::{runtime_config, server_state, shared_engine};
use super::collect;
use crate::test_support::StateUnderTest;
use crate::{
CallerIdentity, NamespaceResolver, ServerState, StaticScheduleNamespaces,
StaticWorkflowNamespaces, config::NamespaceMode,
};
type TestResult = Result<(), Box<dyn std::error::Error>>;
const NAMESPACE: &str = "default";
const OTHER_NAMESPACE: &str = "tenant-b";
struct Fleet {
state: StateUnderTest,
engine: Arc<aion::Engine>,
ownership: Arc<StaticWorkflowNamespaces>,
}
async fn fleet() -> Result<Fleet, Box<dyn std::error::Error>> {
let (engine, store, visibility) = shared_engine().await?;
std::hint::black_box((store, visibility));
let engine_handle = engine.handle();
let ownership = Arc::new(StaticWorkflowNamespaces::default());
let resolver = NamespaceResolver::from_parts(
NamespaceMode::SharedEngine,
Some(engine.handle()),
ownership.clone(),
Arc::new(StaticScheduleNamespaces::default()),
);
let mut runtime = runtime_config();
runtime.auth.enabled = false;
let state = server_state(engine, resolver, runtime).await?;
Ok(Fleet {
state,
engine: engine_handle,
ownership,
})
}
fn observed(minute: i64) -> DateTime<Utc> {
DateTime::UNIX_EPOCH + TimeDelta::minutes(minute)
}
fn entry(reason: &str, minute: i64) -> UnrecoverableRun {
UnrecoverableRun {
workflow_type: String::from("rig$deadbeef"),
reason: reason.to_owned(),
observed_at: observed(minute),
}
}
fn degrade(
fleet: &Fleet,
workflow_id: &WorkflowId,
namespace: Option<&str>,
entry: UnrecoverableRun,
) -> TestResult {
if let Some(namespace) = namespace {
fleet.ownership.record(workflow_id.clone(), namespace)?;
}
fleet
.engine
.registry()
.unrecoverable()
.record(workflow_id.clone(), entry)?;
Ok(())
}
async fn read(
state: ServerState,
) -> Result<(StatusCode, serde_json::Value), Box<dyn std::error::Error>> {
let response = workflow_router(state)
.oneshot(
Request::builder()
.method("GET")
.uri("/workflows/unrecoverable")
.body(Body::empty())?,
)
.await?;
let status = response.status();
let bytes = axum::body::to_bytes(response.into_body(), usize::MAX).await?;
Ok((status, serde_json::from_slice(&bytes)?))
}
#[tokio::test]
async fn a_healthy_engine_reports_no_unrecoverable_runs() -> TestResult {
let fleet = fleet().await?;
let (status, body) = read(fleet.state.clone()).await?;
assert_eq!(status, StatusCode::OK);
assert_eq!(body, serde_json::json!([]));
Ok(())
}
#[tokio::test]
async fn a_degraded_run_is_readable_with_its_id_reason_and_remedy() -> TestResult {
let fleet = fleet().await?;
let workflow_id = WorkflowId::new_v4();
degrade(
&fleet,
&workflow_id,
Some(NAMESPACE),
entry("pinned version 13958627 is not loaded", 0),
)?;
let (status, body) = read(fleet.state.clone()).await?;
assert_eq!(status, StatusCode::OK);
let rows = body.as_array().ok_or("response is not an array")?;
assert_eq!(rows.len(), 1, "{body}");
let row = &rows[0];
assert_eq!(row["workflow_id"], workflow_id.to_string());
assert_eq!(row["namespace"], NAMESPACE);
assert_eq!(row["workflow_type"], "rig$deadbeef");
assert_eq!(row["reason"], "pinned version 13958627 is not loaded");
assert_eq!(row["observed_at"], observed(0).to_rfc3339());
let remedy = row["remedy"].as_str().ok_or("remedy is not a string")?;
assert!(remedy.contains("redeploy"), "{remedy}");
assert!(remedy.contains("cancel"), "{remedy}");
Ok(())
}
#[tokio::test]
async fn a_recovered_run_leaves_the_read() -> TestResult {
let fleet = fleet().await?;
let workflow_id = WorkflowId::new_v4();
degrade(
&fleet,
&workflow_id,
Some(NAMESPACE),
entry("pinned version is not loaded", 0),
)?;
let (_, before) = read(fleet.state.clone()).await?;
assert_eq!(before.as_array().map(Vec::len), Some(1), "{before}");
assert!(
fleet
.engine
.registry()
.unrecoverable()
.clear(&workflow_id)?,
"the fixture must really have removed an entry"
);
let (status, after) = read(fleet.state.clone()).await?;
assert_eq!(status, StatusCode::OK);
assert_eq!(after, serde_json::json!([]));
Ok(())
}
#[tokio::test]
async fn rows_are_ordered_by_when_the_failure_was_observed() -> TestResult {
let fleet = fleet().await?;
let newest = WorkflowId::new_v4();
let oldest = WorkflowId::new_v4();
degrade(&fleet, &newest, Some(NAMESPACE), entry("second", 9))?;
degrade(&fleet, &oldest, Some(NAMESPACE), entry("first", 1))?;
let (_, body) = read(fleet.state.clone()).await?;
let rows = body.as_array().ok_or("response is not an array")?;
assert_eq!(rows.len(), 2, "{body}");
assert_eq!(rows[0]["workflow_id"], oldest.to_string(), "{body}");
assert_eq!(rows[1]["workflow_id"], newest.to_string(), "{body}");
Ok(())
}
#[tokio::test]
async fn a_caller_never_learns_of_a_degraded_run_it_cannot_access() -> TestResult {
let fleet = fleet().await?;
let mine = WorkflowId::new_v4();
let theirs = WorkflowId::new_v4();
degrade(&fleet, &mine, Some(NAMESPACE), entry("mine", 0))?;
degrade(&fleet, &theirs, Some(OTHER_NAMESPACE), entry("theirs", 1))?;
let caller = CallerIdentity::new("tenant-a-subject", [String::from(NAMESPACE)]);
let rows = collect(&fleet.state, &caller).await?;
let ids: Vec<&str> = rows.iter().map(|row| row.workflow_id.as_str()).collect();
assert_eq!(
ids,
vec![mine.to_string().as_str()],
"a tenant must see its own degraded run and only that"
);
let operator = CallerIdentity::operator("operator");
assert_eq!(
collect(&fleet.state, &operator).await?.len(),
2,
"the fixture must hold both runs, or the tenant assertion is vacuous"
);
Ok(())
}
#[tokio::test]
async fn an_unattributed_run_reaches_the_operator_and_stops_there() -> TestResult {
let fleet = fleet().await?;
let orphan = WorkflowId::new_v4();
degrade(&fleet, &orphan, None, entry("no namespace recorded", 0))?;
let operator = CallerIdentity::operator("operator");
let seen = collect(&fleet.state, &operator).await?;
assert_eq!(seen.len(), 1, "the operator must see an unattributed run");
assert_eq!(seen[0].workflow_id, orphan.to_string());
assert!(
seen[0].namespace.is_none(),
"an unattributed run must report a null namespace rather than inventing one"
);
let tenant = CallerIdentity::new("tenant-a-subject", [String::from(NAMESPACE)]);
assert!(
collect(&fleet.state, &tenant).await?.is_empty(),
"an enumerated caller must not be told an unattributed run exists"
);
Ok(())
}