#![allow(dead_code)]
mod admin;
mod bench;
mod cache;
mod cmf_runtime;
mod config;
mod error;
mod import;
mod logbuf;
mod media;
mod model;
mod pipeline;
mod promotion;
mod protocols;
mod providers;
mod registry;
mod router_client;
mod routing;
mod secrets;
mod state;
mod stats;
mod sysmon;
use axum::{routing::get, Router};
use clap::Parser;
use config::Config;
use state::SharedState;
#[derive(Parser)]
#[command(
name = "cortiq-gateway",
version,
about = "Universal LLM gateway with intelligent routing"
)]
struct Args {
#[arg(long)]
config: Option<String>,
#[arg(long)]
admin_token: Option<String>,
#[arg(long)]
no_browser: bool,
#[arg(long)]
no_port_search: bool,
}
fn resolve_admin_token(cfg: &Config, args: &Args) -> String {
if let Some(env) = &cfg.admin.token_env {
if let Ok(v) = std::env::var(env) {
if !v.is_empty() {
return v;
}
}
}
if let Some(t) = &args.admin_token {
if !t.is_empty() {
return t.clone();
}
}
if let Some(t) = &cfg.admin.token {
if !t.is_empty() {
return t.clone();
}
}
String::new()
}
fn is_docker() -> bool {
std::path::Path::new("/.dockerenv").exists()
|| std::fs::read_to_string("/proc/1/cgroup")
.map(|s| s.contains("docker") || s.contains("kubepods"))
.unwrap_or(false)
}
fn should_open_browser(args: &Args) -> bool {
if args.no_browser {
return false;
}
if std::env::var("CORTIQ_NO_BROWSER")
.map(|v| !v.is_empty() && v != "0")
.unwrap_or(false)
{
return false;
}
if std::env::var("CI").is_ok() || std::env::var("NO_BROWSER").is_ok() {
return false;
}
if is_docker() {
return false;
}
true
}
fn parse_listen(listen: &str) -> (String, u16) {
if let Some((host, port_s)) = listen.rsplit_once(':') {
if let Ok(p) = port_s.parse::<u16>() {
return (host.to_string(), p);
}
}
(listen.to_string(), 9000)
}
fn is_port_free(host: &str, port: u16) -> bool {
let addr = format!("{host}:{port}");
if std::net::TcpListener::bind(&addr).is_err() {
return false;
}
use std::net::ToSocketAddrs;
let probe_host = if host == "0.0.0.0" { "127.0.0.1" } else { host };
if let Ok(mut addrs) = (probe_host, port).to_socket_addrs() {
if let Some(a) = addrs.next() {
if std::net::TcpStream::connect_timeout(&a, std::time::Duration::from_millis(200))
.is_ok()
{
return false;
}
}
}
true
}
fn resolve_listen(cfg: &mut Config, no_search: bool) -> String {
if no_search {
return cfg.listen.clone();
}
let (host, port) = parse_listen(&cfg.listen);
let clean_host = if host.is_empty() {
"0.0.0.0".to_string()
} else {
host
};
if is_port_free(&clean_host, port) {
return cfg.listen.clone();
}
for p in (port + 1)..=port.saturating_add(50) {
if is_port_free(&clean_host, p) {
let new_listen = format!("{clean_host}:{p}");
eprintln!("port {port} busy — switching to free port {p} ({new_listen})");
cfg.listen = new_listen.clone();
return new_listen;
}
}
if !is_port_free(&clean_host, port) {
eprintln!(
"warning: port {port} busy and no free port found in next 50, will try to bind anyway"
);
}
cfg.listen.clone()
}
fn setup_marker_path(config_path: &str) -> std::path::PathBuf {
std::path::Path::new(config_path)
.parent()
.unwrap_or_else(|| std::path::Path::new("."))
.join(".setup_done")
}
fn needs_first_run_wizard(cfg: &Config, config_path: &str) -> bool {
let marker = setup_marker_path(config_path);
if marker.exists() {
return false;
}
if !cfg.models.is_empty() || !cfg.cmf.effective_servers().is_empty() {
let _ = std::fs::write(&marker, "done");
return false;
}
true
}
#[tokio::main]
async fn main() -> anyhow::Result<()> {
let args = Args::parse();
let config_path = args
.config
.clone()
.unwrap_or_else(config::default_config_path);
let mut cfg = Config::load_or_create(&config_path)?;
let original_listen = cfg.listen.clone();
let listen = resolve_listen(&mut cfg, args.no_port_search);
if listen != original_listen {
let _ = cfg.save(&config_path);
}
{
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
let filter = tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| cfg.log.level.clone().into());
tracing_subscriber::registry()
.with(filter)
.with(tracing_subscriber::fmt::layer())
.with(
tracing_subscriber::fmt::layer()
.with_ansi(false)
.with_writer(logbuf::RingMakeWriter),
)
.init();
}
let admin_enabled = cfg.admin.enabled;
let admin_token = resolve_admin_token(&cfg, &args);
let first_run = needs_first_run_wizard(&cfg, &config_path);
let state = SharedState::build(cfg, config_path.clone())?;
let host_public = listen.starts_with("0.0.0.0") || listen.starts_with("[::]");
if admin_enabled && admin_token.is_empty() && host_public && !is_docker() {
state
.admin_open_public
.store(true, std::sync::atomic::Ordering::Relaxed);
tracing::warn!(
"admin console is OPEN (no token) and listening on {listen} — anyone on your network can manage this gateway; set a token via --admin-token or [admin].token if this machine is reachable from outside"
);
}
{
let cmf_rt = state.cmf.clone();
let cmf_cfg = state.live().cfg.cmf.clone();
tokio::spawn(async move { cmf_runtime::manage(cmf_rt, cmf_cfg).await });
}
let mut app = Router::new()
.route("/healthz", get(|| async { "ok" }))
.route("/health", get(|| async { "ok" }))
.route("/readyz", get(|| async { "ready" }))
.route("/metrics", get(admin::metrics))
.merge(protocols::build_router(state.clone()));
if admin_enabled {
app = app
.merge(admin::api_routes(admin_token.clone()))
.fallback(admin::assets::fallback);
}
let cors = tower_http::cors::CorsLayer::new()
.allow_origin(tower_http::cors::Any)
.allow_methods(tower_http::cors::Any)
.allow_headers(tower_http::cors::Any);
let app =
app.with_state(state)
.layer(cors)
.layer(tower_http::limit::RequestBodyLimitLayer::new(
32 * 1024 * 1024,
));
let listener = {
let mut attempt = 0u32;
loop {
match tokio::net::TcpListener::bind(&listen).await {
Ok(l) => break l,
Err(e) => {
let retrying = std::env::var("CORTIQ_BIND_RETRY").is_ok();
if retrying && attempt < 50 {
attempt += 1;
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
continue;
}
return Err(e.into());
}
}
}
};
tracing::info!("cortiq-gateway listening on {listen}");
if admin_enabled {
let host_for_url = if listen.starts_with("0.0.0.0:") {
listen.replacen("0.0.0.0", "127.0.0.1", 1)
} else {
listen.clone()
};
let admin_path = if first_run {
"/admin#/onboarding"
} else {
"/admin"
};
let url = if admin_token.is_empty() {
format!("http://{host_for_url}{admin_path}")
} else {
format!("http://{host_for_url}{admin_path}?token={admin_token}")
};
if admin_token.is_empty() {
tracing::info!("admin console (open, no token): http://{host_for_url}/admin");
} else {
tracing::info!("admin console: http://{host_for_url}/admin");
}
if first_run {
tracing::info!("first run — opening onboarding wizard at {url}");
}
if admin_enabled && should_open_browser(&args) {
let browser_url = url.clone();
tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_millis(600)).await;
#[cfg(feature = "open-browser")]
{
if let Err(e) = open::that(&browser_url) {
tracing::warn!("failed to open browser: {e}");
} else {
tracing::info!("opened browser at {browser_url}");
}
}
#[cfg(not(feature = "open-browser"))]
{
tracing::info!("open browser: {browser_url} (feature disabled)");
}
});
} else if first_run {
tracing::info!("open this URL to finish setup: {url}");
}
}
axum::serve(listener, app)
.with_graceful_shutdown(async {
let _ = tokio::signal::ctrl_c().await;
tracing::info!("shutdown signal received — draining connections");
})
.await?;
Ok(())
}