use std::time::Duration;
use axum::{
extract::{Path, State},
http::HeaderMap,
Json,
};
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use uuid::Uuid;
use crate::{
error::{ApiError, ApiResult},
middleware::{resolve_org_context, AuthUser},
models::HostedMock,
AppState,
};
const RUNTIME_HTTP_PORT: u16 = 3000;
const PROXY_TIMEOUT: Duration = Duration::from_secs(3);
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct TimeTravelStatusResponse {
pub enabled: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub current_time: Option<String>,
pub scale_factor: f64,
pub real_time: String,
}
impl TimeTravelStatusResponse {
fn disabled_placeholder() -> Self {
Self {
enabled: false,
current_time: None,
scale_factor: 1.0,
real_time: chrono::Utc::now().to_rfc3339(),
}
}
}
#[derive(Debug, Serialize)]
pub struct TimeTravelEnvelope<T: Serialize> {
pub runtime_state: &'static str,
pub data: T,
}
impl<T: Serialize> TimeTravelEnvelope<T> {
fn live(data: T) -> Self {
Self {
runtime_state: "live",
data,
}
}
fn unreachable(data: T) -> Self {
Self {
runtime_state: "unreachable",
data,
}
}
}
#[derive(Debug, Deserialize, Serialize)]
pub struct EnableRequest {
#[serde(skip_serializing_if = "Option::is_none")]
pub time: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub scale: Option<f64>,
}
#[derive(Debug, Deserialize, Serialize)]
pub struct AdvanceRequest {
pub duration: String,
}
#[derive(Debug, Deserialize, Serialize)]
pub struct SetTimeRequest {
pub time: String,
}
#[derive(Debug, Deserialize, Serialize)]
pub struct SetScaleRequest {
pub scale: f64,
}
pub async fn get_status(
State(state): State<AppState>,
AuthUser(user_id): AuthUser,
Path(deployment_id): Path<Uuid>,
headers: HeaderMap,
) -> ApiResult<Json<TimeTravelEnvelope<TimeTravelStatusResponse>>> {
let deployment = authorize_deployment(&state, user_id, &headers, deployment_id).await?;
let url = format!("{}/__mockforge/time-travel/status", runtime_base_url(&deployment));
Ok(Json(match proxy_get::<TimeTravelStatusResponse>(&url).await {
Ok(data) => TimeTravelEnvelope::live(data),
Err(err) => {
tracing::warn!(%deployment_id, error = %err, "time-travel proxy GET status failed");
TimeTravelEnvelope::unreachable(TimeTravelStatusResponse::disabled_placeholder())
}
}))
}
pub async fn enable(
State(state): State<AppState>,
AuthUser(user_id): AuthUser,
Path(deployment_id): Path<Uuid>,
headers: HeaderMap,
Json(body): Json<EnableRequest>,
) -> ApiResult<Json<Value>> {
let deployment = authorize_deployment(&state, user_id, &headers, deployment_id).await?;
let url = format!("{}/__mockforge/time-travel/enable", runtime_base_url(&deployment));
Ok(Json(proxy_post_json(&url, &body, deployment_id, "enable").await))
}
pub async fn disable(
State(state): State<AppState>,
AuthUser(user_id): AuthUser,
Path(deployment_id): Path<Uuid>,
headers: HeaderMap,
) -> ApiResult<Json<Value>> {
let deployment = authorize_deployment(&state, user_id, &headers, deployment_id).await?;
let url = format!("{}/__mockforge/time-travel/disable", runtime_base_url(&deployment));
Ok(Json(proxy_post_json(&url, &json!({}), deployment_id, "disable").await))
}
pub async fn advance(
State(state): State<AppState>,
AuthUser(user_id): AuthUser,
Path(deployment_id): Path<Uuid>,
headers: HeaderMap,
Json(body): Json<AdvanceRequest>,
) -> ApiResult<Json<Value>> {
let deployment = authorize_deployment(&state, user_id, &headers, deployment_id).await?;
let url = format!("{}/__mockforge/time-travel/advance", runtime_base_url(&deployment));
Ok(Json(proxy_post_json(&url, &body, deployment_id, "advance").await))
}
pub async fn set_time(
State(state): State<AppState>,
AuthUser(user_id): AuthUser,
Path(deployment_id): Path<Uuid>,
headers: HeaderMap,
Json(body): Json<SetTimeRequest>,
) -> ApiResult<Json<Value>> {
let deployment = authorize_deployment(&state, user_id, &headers, deployment_id).await?;
let url = format!("{}/__mockforge/time-travel/set", runtime_base_url(&deployment));
Ok(Json(proxy_post_json(&url, &body, deployment_id, "set").await))
}
pub async fn set_scale(
State(state): State<AppState>,
AuthUser(user_id): AuthUser,
Path(deployment_id): Path<Uuid>,
headers: HeaderMap,
Json(body): Json<SetScaleRequest>,
) -> ApiResult<Json<Value>> {
let deployment = authorize_deployment(&state, user_id, &headers, deployment_id).await?;
let url = format!("{}/__mockforge/time-travel/scale", runtime_base_url(&deployment));
Ok(Json(proxy_post_json(&url, &body, deployment_id, "scale").await))
}
pub async fn reset(
State(state): State<AppState>,
AuthUser(user_id): AuthUser,
Path(deployment_id): Path<Uuid>,
headers: HeaderMap,
) -> ApiResult<Json<Value>> {
let deployment = authorize_deployment(&state, user_id, &headers, deployment_id).await?;
let url = format!("{}/__mockforge/time-travel/reset", runtime_base_url(&deployment));
Ok(Json(proxy_post_json(&url, &json!({}), deployment_id, "reset").await))
}
fn runtime_base_url(deployment: &HostedMock) -> String {
format!("http://{}.internal:{RUNTIME_HTTP_PORT}", deployment.fly_app_name())
}
async fn proxy_get<T: for<'de> Deserialize<'de>>(url: &str) -> reqwest::Result<T> {
reqwest::Client::builder()
.timeout(PROXY_TIMEOUT)
.build()?
.get(url)
.send()
.await?
.error_for_status()?
.json::<T>()
.await
}
async fn proxy_post_json<B: Serialize>(
url: &str,
body: &B,
deployment_id: Uuid,
op: &'static str,
) -> Value {
let client = match reqwest::Client::builder().timeout(PROXY_TIMEOUT).build() {
Ok(c) => c,
Err(err) => {
tracing::warn!(%deployment_id, op, error = %err, "reqwest client build failed");
return unreachable_post_body(err.to_string());
}
};
let resp = match client.post(url).json(body).send().await {
Ok(r) => r,
Err(err) => {
tracing::warn!(%deployment_id, op, error = %err, "time-travel POST failed");
return unreachable_post_body(err.to_string());
}
};
if !resp.status().is_success() {
let status = resp.status();
let text = resp.text().await.unwrap_or_default();
tracing::warn!(%deployment_id, op, %status, body = %text, "time-travel POST non-2xx");
return unreachable_post_body(format!("upstream {status}"));
}
let upstream: Value = resp.json().await.unwrap_or(Value::Null);
let mut body = json!({
"accepted": true,
"runtime_state": "live",
});
if let Value::Object(ref mut map) = body {
if !upstream.is_null() {
map.insert("upstream".into(), upstream);
}
}
body
}
fn unreachable_post_body(reason: String) -> Value {
json!({
"accepted": false,
"runtime_state": "unreachable",
"reason": reason,
})
}
async fn authorize_deployment(
state: &AppState,
user_id: Uuid,
headers: &HeaderMap,
deployment_id: Uuid,
) -> ApiResult<HostedMock> {
let deployment = HostedMock::find_by_id(state.db.pool(), deployment_id)
.await?
.ok_or_else(|| ApiError::InvalidRequest("Deployment not found".into()))?;
let ctx = resolve_org_context(state, user_id, headers, None)
.await
.map_err(|_| ApiError::InvalidRequest("Organization not found".into()))?;
if ctx.org_id != deployment.org_id {
return Err(ApiError::InvalidRequest("Deployment not found".into()));
}
Ok(deployment)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn envelope_live_round_trips_status() {
let env = TimeTravelEnvelope::live(TimeTravelStatusResponse {
enabled: true,
current_time: Some("2030-01-01T00:00:00Z".into()),
scale_factor: 60.0,
real_time: "2026-05-16T05:00:00Z".into(),
});
let body = serde_json::to_value(&env).unwrap();
assert_eq!(body["runtime_state"], "live");
assert_eq!(body["data"]["enabled"], true);
assert_eq!(body["data"]["current_time"], "2030-01-01T00:00:00Z");
assert_eq!(body["data"]["scale_factor"], 60.0);
}
#[test]
fn envelope_unreachable_is_disabled_placeholder() {
let env = TimeTravelEnvelope::unreachable(TimeTravelStatusResponse::disabled_placeholder());
let body = serde_json::to_value(&env).unwrap();
assert_eq!(body["runtime_state"], "unreachable");
assert_eq!(body["data"]["enabled"], false);
assert_eq!(body["data"]["scale_factor"], 1.0);
assert!(
body["data"]["current_time"].is_null()
|| !body["data"].as_object().unwrap().contains_key("current_time")
);
}
#[test]
fn unreachable_post_body_shape() {
let body = unreachable_post_body("connection refused".into());
assert_eq!(body["accepted"], false);
assert_eq!(body["runtime_state"], "unreachable");
assert_eq!(body["reason"], "connection refused");
}
#[test]
fn enable_request_serializes_minimal() {
let body = serde_json::to_value(EnableRequest {
time: None,
scale: None,
})
.unwrap();
assert_eq!(body, serde_json::json!({}));
}
#[test]
fn enable_request_serializes_partial() {
let body = serde_json::to_value(EnableRequest {
time: None,
scale: Some(60.0),
})
.unwrap();
assert_eq!(body["scale"], 60.0);
assert!(body.get("time").is_none());
}
}