#![cfg(feature = "prometheus")]
use std::collections::HashMap;
use std::sync::Arc;
use std::time::{Duration, Instant};
use axum::serve;
use chrono::{TimeDelta, Utc};
use ironflow_api::routes::{RouterConfig, create_router};
use ironflow_api::state::AppState;
use ironflow_auth::jwt::JwtConfig;
use ironflow_core::metric_names::WORKER_QUEUE_DEPTH;
use ironflow_core::providers::claude::ClaudeCodeProvider;
use ironflow_engine::engine::{Engine, ExecutionMode};
use ironflow_engine::notify::Event;
use ironflow_store::memory::InMemoryStore;
use ironflow_store::models::{NewRun, RunStatus, TriggerKind};
use ironflow_store::store::RunStore;
use ironflow_worker::WorkerBuilder;
use serde_json::json;
use tokio::net::TcpListener;
use tokio::spawn;
use tokio::sync::broadcast;
use tokio::time::{sleep, timeout};
const TEST_TIMEOUT: Duration = Duration::from_secs(30);
fn future_run(workflow: &str) -> NewRun {
NewRun {
created_by: None,
workflow_name: workflow.to_string(),
trigger: TriggerKind::Manual,
payload: json!({}),
max_retries: 0,
handler_version: None,
labels: HashMap::new(),
scheduled_at: Some(Utc::now() + TimeDelta::hours(1)),
idempotency_key: None,
concurrency_key: None,
max_cost_usd: None,
}
}
fn gauge_value(exposition: &str, metric: &str) -> Option<f64> {
exposition
.lines()
.filter(|line| !line.starts_with('#'))
.find_map(|line| {
let (name, value) = line.split_once(' ')?;
(name == metric)
.then(|| value.trim().parse().ok())
.flatten()
})
}
#[tokio::test]
async fn api_mode_worker_publishes_pending_run_count_as_queue_depth() {
timeout(TEST_TIMEOUT, async {
let store = Arc::new(InMemoryStore::new());
for _ in 0..3 {
store
.create_run(future_run("queued"))
.await
.expect("create pending run");
}
let running = store
.create_run(future_run("queued"))
.await
.expect("create run")
.into_run();
store
.update_run_status(running.id, RunStatus::Running)
.await
.expect("move run to running");
let engine = Engine::new(store.clone(), Arc::new(ClaudeCodeProvider::new()))
.with_execution_mode(ExecutionMode::Workers);
let jwt_config = Arc::new(JwtConfig {
secret: "test-secret-for-queue-depth".to_string(),
access_token_ttl_secs: 900,
refresh_token_ttl_secs: 604800,
cookie_domain: None,
cookie_secure: false,
});
let (event_sender, _) = broadcast::channel::<Event>(16);
let state = AppState::new(
store.clone(),
Arc::new(engine),
jwt_config,
"test-worker-token".to_string(),
event_sender,
);
let prometheus = state.prometheus_handle.clone();
let config = RouterConfig {
rate_limit_auth: None,
rate_limit_general: None,
..RouterConfig::default()
};
let router = create_router(state, config);
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("local addr");
spawn(async move {
serve(listener, router).await.expect("serve");
});
let worker = WorkerBuilder::new(&format!("http://{addr}"), "test-worker-token")
.provider(Arc::new(ClaudeCodeProvider::new()))
.worker_id("worker-queue-depth")
.concurrency(1)
.poll_interval(Duration::from_millis(20))
.build()
.expect("build worker");
let handle = spawn(async move {
if let Err(e) = worker.run().await {
eprintln!("worker exited with error: {e:?}");
}
});
let deadline = Instant::now() + Duration::from_secs(8);
let mut depth = None;
while Instant::now() < deadline {
depth = gauge_value(&prometheus.render(), WORKER_QUEUE_DEPTH);
if depth == Some(3.0) {
break;
}
sleep(Duration::from_millis(20)).await;
}
handle.abort();
assert_eq!(
depth,
Some(3.0),
"expected {WORKER_QUEUE_DEPTH} 3, exposition was:\n{}",
prometheus.render()
);
})
.await
.expect("test timed out");
}