use optionchain_simulator::api::start_server;
use optionchain_simulator::infrastructure::{
ClickHouseSnapshotRepository, DEFAULT_PRICING_GATE_KEY, DependencyProbe, MetricsCollector,
MongoDbProbe, Readiness, RedisClient, RedisConfig, RedisPricingGate, RedisProbe, ServerConfig,
SimulationV2Config, WarehouseProbe, init_mongodb, resolve_log_level_from_env,
};
use optionchain_simulator::session::{
DEFAULT_TAPE_KEY_PREFIX, InRedisSessionStore, InRedisSimulationStore, RedisTapeCache,
SessionManager, SimulationManager,
};
use optionchain_simulator::utils::admission::{configured_jobs, install_shared_gate};
use optionstratlib::utils::setup_logger_with_level;
use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use tracing::{info, warn};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let log_level = resolve_log_level_from_env();
setup_logger_with_level(log_level.level.as_str());
if let Some(rejected) = &log_level.rejected {
warn!(
value = %rejected,
default = %log_level.level,
"unrecognised LOGLEVEL; using the default"
);
}
info!(level = %log_level.level, "Log level resolved from LOGLEVEL");
let server = ServerConfig::from_env()?;
let redis_config = RedisConfig::default();
info!("Connecting to Redis at {}", redis_config);
let redis_client = Arc::new(RedisClient::new(redis_config).await?);
let redis_client_v2 = Arc::clone(&redis_client);
let redis_client_probe = Arc::clone(&redis_client);
let store = Arc::new(InRedisSessionStore::new(
redis_client,
Some("optionchain:session:".to_string()), Some(3600), ));
let metrics_collector = Arc::new(MetricsCollector::new()?);
let mongodb_repository = init_mongodb().await?;
let session_manager = Arc::new(SessionManager::new(store.clone()));
let v2_config = SimulationV2Config::from_env()?;
let simulation_store = Arc::new(InRedisSimulationStore::new(
Arc::clone(&redis_client_v2),
None, Some(v2_config.retention_secs()),
));
let pricing_gate = Arc::new(RedisPricingGate::new(
Arc::clone(&redis_client_v2),
DEFAULT_PRICING_GATE_KEY,
configured_jobs(),
));
install_shared_gate(pricing_gate)?;
let shared_tapes = Arc::new(RedisTapeCache::new(
Arc::clone(&redis_client_v2),
DEFAULT_TAPE_KEY_PREFIX,
Duration::from_secs(v2_config.retention_secs()),
v2_config.max_cached_tapes,
));
let mut simulation_manager =
SimulationManager::new(simulation_store, v2_config).with_shared_tapes(shared_tapes);
match ClickHouseSnapshotRepository::from_env()? {
Some(warehouse) => {
warehouse.ensure_schema().await?;
info!("v2 snapshot persistence is enabled");
simulation_manager = simulation_manager
.with_snapshot_metrics(Arc::clone(&metrics_collector))
.with_warehouse(Arc::new(warehouse));
}
None => {
info!("v2 snapshot persistence is disabled; snapshots are served from replay only");
}
}
let simulation_manager = Arc::new(simulation_manager);
let sweeper = Arc::clone(&simulation_manager);
let sweeper_metrics = Arc::clone(&metrics_collector);
tokio::spawn(async move {
let mut ticker = tokio::time::interval(v2_config.cleanup_interval);
loop {
ticker.tick().await;
match sweeper.cleanup().await {
Ok(expired) => {
sweeper_metrics.record_v2_simulations_expired(expired.len());
sweeper_metrics.set_v2_cache_sizes(
sweeper.cached_tapes() as i64,
sweeper.cached_snapshots() as i64,
);
}
Err(error) => warn!(%error, "the v2 retention sweep failed"),
}
}
});
let mut probes: Vec<Arc<dyn DependencyProbe>> = vec![
Arc::new(RedisProbe::new(redis_client_probe)),
Arc::new(MongoDbProbe::new(Arc::clone(&mongodb_repository))),
];
if let Some(warehouse) = simulation_manager.warehouse() {
probes.push(Arc::new(WarehouseProbe::new(warehouse)));
}
let readiness = Readiness::new(probes);
let listen_on = server.address;
let port = server.port;
info!(
"Starting HTTP server at http://{}",
SocketAddr::new(listen_on.ip(), port)
);
if listen_on.is_public() {
warn!(
address = %listen_on,
"Listening beyond loopback; the service has no authentication"
);
}
match start_server(
session_manager,
simulation_manager,
metrics_collector,
mongodb_repository,
listen_on,
port,
readiness,
)
.await
{
Ok(_) => Ok(()),
Err(e) => Err(e.to_string().into()),
}
}