use std::sync::Mutex;
use std::time::{Duration, Instant};
use crate::config::Configs;
use crate::telemetry::{self, CliTrackEvent};
static STAGES: Mutex<Vec<StageTiming>> = Mutex::new(Vec::new());
struct StageTiming {
stage: String,
duration_ms: u64,
success: bool,
}
fn timing_to_stderr() -> bool {
std::env::var("RAILWAY_STAGE_TIMING").is_ok_and(|v| v == "1" || v == "true")
}
pub fn record_stage(stage: &str, duration: Duration, success: bool) {
record(stage, duration, success);
}
static DETACHED: Mutex<Vec<tokio::task::JoinHandle<()>>> = Mutex::new(Vec::new());
pub fn spawn_detached(fut: impl std::future::Future<Output = ()> + Send + 'static) {
let handle = tokio::spawn(fut);
if let Ok(mut detached) = DETACHED.lock() {
detached.push(handle);
}
}
pub async fn drain_detached(limit: Duration) {
let handles = match DETACHED.lock() {
Ok(mut guard) => std::mem::take(&mut *guard),
Err(_) => return,
};
let deadline = Instant::now() + limit;
for handle in handles {
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return;
}
let _ = tokio::time::timeout(remaining, handle).await;
}
}
pub fn timing_diagnostics() -> bool {
timing_to_stderr()
}
fn record(stage: &str, duration: Duration, success: bool) {
let timing = StageTiming {
stage: stage.to_string(),
duration_ms: duration.as_millis() as u64,
success,
};
if let Ok(mut stages) = STAGES.lock() {
stages.push(timing);
}
}
pub async fn timed_for<T>(
command: &str,
stage: &str,
fut: impl std::future::Future<Output = anyhow::Result<T>>,
) -> anyhow::Result<T> {
let started = Instant::now();
let result = fut.await;
let elapsed = started.elapsed();
record(stage, elapsed, result.is_ok());
if let Err(ref e) = result {
report_failure_timed(command, stage, &format!("{e}"), elapsed.as_millis() as u64).await;
}
result
}
pub fn flush_stages(command: &'static str) {
let stages = match STAGES.lock() {
Ok(mut guard) => std::mem::take(&mut *guard),
Err(_) => return,
};
if stages.is_empty() {
return;
}
if timing_to_stderr() {
let total: u64 = stages.iter().map(|s| s.duration_ms).sum();
let mut line = format!("[{command} timing] total_tracked={total}ms");
for s in &stages {
let mark = if s.success { "" } else { "!" };
line.push_str(&format!(" {}{}={}ms", s.stage, mark, s.duration_ms));
}
eprintln!("{line}");
}
for s in stages {
if !s.success {
continue;
}
spawn_detached(telemetry::send(CliTrackEvent {
command: command.to_string(),
sub_command: Some(format!("stage_{}", s.stage)),
success: true,
error_message: None,
duration_ms: s.duration_ms,
cli_version: env!("CARGO_PKG_VERSION"),
os: std::env::consts::OS,
arch: std::env::consts::ARCH,
is_ci: Configs::env_is_ci(),
}));
}
}
pub async fn report_failure(stage: &str, message: &str) {
report_failure_for("ssh", stage, message).await;
}
pub async fn report_failure_for(command: &str, stage: &str, message: &str) {
report_failure_timed(command, stage, message, 0).await;
}
pub async fn report_failure_timed(command: &str, stage: &str, message: &str, duration_ms: u64) {
let truncated = telemetry::truncate_message(message);
telemetry::send(CliTrackEvent {
command: command.to_string(),
sub_command: Some(format!("stage_{stage}_failed")),
success: false,
error_message: Some(truncated),
duration_ms,
cli_version: env!("CARGO_PKG_VERSION"),
os: std::env::consts::OS,
arch: std::env::consts::ARCH,
is_ci: Configs::env_is_ci(),
})
.await;
}
pub async fn track<T>(stage: &str, result: anyhow::Result<T>) -> anyhow::Result<T> {
track_for("ssh", stage, result).await
}
pub async fn track_for<T>(
command: &str,
stage: &str,
result: anyhow::Result<T>,
) -> anyhow::Result<T> {
if let Err(ref e) = result {
report_failure_for(command, stage, &format!("{e}")).await;
}
result
}