use async_nats::{Client, ConnectOptions, Subscriber};
use chrono::{DateTime, Utc};
use serde::Deserialize;
use std::time::Duration;
use tracing::{debug, info};
use crate::fly_logs::LogEntry;
const DEFAULT_NATS_URL: &str = "nats://[fdaa::3]:4223";
const DEFAULT_CONNECT_TIMEOUT_MS: u64 = 1500;
#[derive(Debug, Clone)]
pub struct FlyNatsConfig {
pub url: String,
pub org_slug: String,
pub token: String,
pub connect_timeout: Duration,
}
impl FlyNatsConfig {
pub fn from_env() -> Option<Self> {
let org_slug = std::env::var("FLY_ORG_SLUG").ok()?;
let token = std::env::var("FLYIO_API_TOKEN").ok()?;
let url = std::env::var("FLY_NATS_URL").unwrap_or_else(|_| DEFAULT_NATS_URL.to_string());
let connect_timeout_ms = std::env::var("FLY_NATS_CONNECT_TIMEOUT_MS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(DEFAULT_CONNECT_TIMEOUT_MS);
Some(Self {
url,
org_slug,
token,
connect_timeout: Duration::from_millis(connect_timeout_ms),
})
}
}
pub struct FlyLogSubscription {
pub subscriber: Subscriber,
_client: Client,
}
pub async fn subscribe_app(
config: &FlyNatsConfig,
app_name: &str,
) -> Result<FlyLogSubscription, async_nats::Error> {
debug!(app = %app_name, url = %config.url, "connecting to Fly NATS proxy");
let client = ConnectOptions::new()
.user_and_password(config.org_slug.clone(), config.token.clone())
.connection_timeout(config.connect_timeout)
.max_reconnects(Some(0))
.name(format!("mockforge-registry runtime-logs ({app_name})"))
.connect(&config.url)
.await?;
let subject = format!("logs.{app_name}.>");
info!(app = %app_name, subject = %subject, "subscribed to Fly NATS log subject");
let subscriber = client.subscribe(subject).await?;
Ok(FlyLogSubscription {
subscriber,
_client: client,
})
}
pub fn parse_message(payload: &[u8]) -> Option<LogEntry> {
let raw: FlyLogPayload = serde_json::from_slice(payload).ok()?;
let timestamp = raw
.timestamp
.as_deref()
.and_then(|s| DateTime::parse_from_rfc3339(s).ok())
.map(|d| d.with_timezone(&Utc))
.unwrap_or_else(Utc::now);
let level = raw
.log
.as_ref()
.and_then(|l| l.level.clone())
.or(raw.level)
.unwrap_or_else(|| "info".to_string());
let message = raw.message?;
let instance = raw.fly.as_ref().and_then(|f| f.app.as_ref()).and_then(|a| a.instance.clone());
let region = raw.fly.and_then(|f| f.region);
Some(LogEntry {
timestamp,
level,
message,
instance,
region,
})
}
#[derive(Debug, Deserialize)]
struct FlyLogPayload {
#[serde(default)]
timestamp: Option<String>,
#[serde(default)]
level: Option<String>,
#[serde(default)]
message: Option<String>,
#[serde(default)]
log: Option<FlyLogInner>,
#[serde(default)]
fly: Option<FlyEnvelope>,
}
#[derive(Debug, Deserialize)]
struct FlyLogInner {
#[serde(default)]
level: Option<String>,
}
#[derive(Debug, Deserialize)]
struct FlyEnvelope {
#[serde(default)]
region: Option<String>,
#[serde(default)]
app: Option<FlyApp>,
}
#[derive(Debug, Deserialize)]
struct FlyApp {
#[serde(default)]
instance: Option<String>,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn config_requires_org_and_token() {
std::env::remove_var("FLY_ORG_SLUG");
std::env::remove_var("FLYIO_API_TOKEN");
assert!(FlyNatsConfig::from_env().is_none());
std::env::set_var("FLY_ORG_SLUG", "test-org");
assert!(FlyNatsConfig::from_env().is_none(), "org without token should be None");
std::env::set_var("FLYIO_API_TOKEN", "test-token-value");
let cfg = FlyNatsConfig::from_env().expect("both set");
assert_eq!(cfg.org_slug, "test-org");
assert_eq!(cfg.token, "test-token-value");
assert_eq!(cfg.url, DEFAULT_NATS_URL);
std::env::set_var("FLY_NATS_URL", "nats://localhost:4222");
let cfg = FlyNatsConfig::from_env().unwrap();
assert_eq!(cfg.url, "nats://localhost:4222");
std::env::remove_var("FLY_NATS_URL");
std::env::remove_var("FLY_ORG_SLUG");
std::env::remove_var("FLYIO_API_TOKEN");
}
#[test]
fn parses_vector_envelope_form() {
let raw = br#"{
"timestamp": "2026-05-18T10:00:00Z",
"log": { "level": "warn" },
"message": "GET /api/users 500",
"fly": {
"region": "lhr",
"app": { "instance": "machineabc123", "name": "my-app" }
}
}"#;
let entry = parse_message(raw).expect("parses");
assert_eq!(entry.message, "GET /api/users 500");
assert_eq!(entry.level, "warn");
assert_eq!(entry.region.as_deref(), Some("lhr"));
assert_eq!(entry.instance.as_deref(), Some("machineabc123"));
}
#[test]
fn parses_flat_form() {
let raw = br#"{
"timestamp": "2026-05-18T10:00:00Z",
"level": "error",
"message": "boom"
}"#;
let entry = parse_message(raw).expect("parses");
assert_eq!(entry.level, "error");
assert_eq!(entry.message, "boom");
assert!(entry.region.is_none());
}
#[test]
fn drops_message_without_text() {
let raw = br#"{ "timestamp": "2026-05-18T10:00:00Z", "level": "info" }"#;
assert!(parse_message(raw).is_none());
}
#[test]
fn drops_unparsable_payload() {
assert!(parse_message(b"not json").is_none());
}
#[test]
fn defaults_missing_level_to_info() {
let raw = br#"{ "message": "hi", "timestamp": "2026-05-18T10:00:00Z" }"#;
let entry = parse_message(raw).unwrap();
assert_eq!(entry.level, "info");
}
}