#![forbid(unsafe_code)]
use std::io::{self, Stdout};
use std::path::{Path, PathBuf};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use chainview::config::{CliOverrides, Config, ModeSelect};
use chainview::{
AliasCatalog, App, BundleError, ChainFetch, ChainSource, ChainStore, ChainViewApp,
ChainViewError, Command as DataCommand, EventBridge, ExitCause, ExpirySource, GuardTeardown,
Instrument, LiveState, Mode, ReplayState, Resolved, ResourceCeilings, SourceBinding,
Supervisor, TerminalGuard, TokioTask, event_channel, install_panic_hook, run_render_loop,
spawn_bundle_load, spawn_input_reader, spawn_supervised_subscription, spawn_tick_task,
};
use chrono::{DateTime, Utc};
use clap::{Args, Parser, Subcommand};
use optionstratlib::ExpirationDate;
use optionstratlib::chains::chain::OptionChain;
use optionstratlib::prelude::Positive;
use ratatui::Terminal;
use ratatui::backend::CrosstermBackend;
#[derive(Debug, Parser)]
#[command(name = "chainview", version, about, long_about = None)]
struct Cli {
#[command(flatten)]
live: LiveArgs,
#[command(subcommand)]
command: Option<Command>,
}
#[derive(Debug, Subcommand)]
enum Command {
Replay {
dir: PathBuf,
},
}
#[derive(Debug, Args)]
struct LiveArgs {
#[arg(long)]
provider: Option<String>,
#[arg(long)]
underlying: Option<String>,
#[arg(long)]
refresh: Option<String>,
#[arg(long)]
tick: Option<String>,
#[arg(long = "channel-cap")]
channel_cap: Option<i64>,
#[arg(long = "log-file")]
log_file: Option<PathBuf>,
#[arg(long)]
theme: Option<String>,
#[arg(long = "no-color")]
no_color: bool,
#[arg(long)]
endpoint: Option<String>,
}
impl Cli {
fn into_overrides(self) -> CliOverrides {
let mode = match self.command {
Some(Command::Replay { dir }) => ModeSelect::Replay(dir),
None => ModeSelect::Live,
};
let live = self.live;
CliOverrides {
provider: live.provider,
underlying: live.underlying,
refresh_interval: live.refresh,
tick_interval: live.tick,
channel_capacity: live.channel_cap,
log_file: live.log_file,
theme: live.theme,
no_color: live.no_color,
endpoint: live.endpoint,
mode,
}
}
}
fn main() -> Result<(), ChainViewError> {
let _ = dotenvy::dotenv();
let overrides = Cli::parse().into_overrides();
let config = Config::load(overrides)?;
if let ModeSelect::Replay(dir) = &config.mode {
validate_replay_dir(dir)?;
}
install_panic_hook();
match ChainViewApp::builder()
.with_builtins()
.with_config(config)
.resolve()?
{
Resolved::Live {
provider,
source,
config,
} => run_live(provider, source, config),
Resolved::Replay { dir, config } => run_replay(dir, config),
}
}
fn run_live(
provider: std::sync::Arc<dyn chainview::Provider>,
source: SourceBinding,
config: Config,
) -> Result<(), ChainViewError> {
let runtime = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.map_err(|e| ChainViewError::Terminal(format!("tokio runtime: {e}")))?;
let cause = runtime.block_on(compose_and_run_live(provider, source, config));
match cause {
ExitCause::Clean => Ok(()),
ExitCause::TaskPanicked => Err(ChainViewError::Terminal(
"a supervised task panicked; see the log".to_owned(),
)),
ExitCause::Failed(error) => Err(error),
}
}
async fn compose_and_run_live(
provider: std::sync::Arc<dyn chainview::Provider>,
source: SourceBinding,
config: Config,
) -> ExitCause {
let now = now_utc();
let expiration = ExpirationDate::Days(positive_or_one(7.0));
let (fetch, instruments) = match provider.fetch_chain(&config.underlying, &expiration).await {
Ok(fetch) => {
let instruments: Vec<Instrument> = fetch.aliases.instruments().cloned().collect();
(fetch, instruments)
}
Err(_) => (empty_seed(&config.underlying, &source, now), Vec::new()),
};
let expiration_utc = fetch.expiry_source.expiration_utc;
let store = ChainStore::seed(fetch, ChainSource::Merged, config.refresh_interval, now);
let live = LiveState::new(source, store);
let (mut bridge, senders) = EventBridge::new(config.channel_capacity);
let mut app = App::new(Mode::Live(live), config.theme, senders.tx_command.clone())
.with_no_color(config.no_color);
let guard = match TerminalGuard::new() {
Ok(guard) => guard,
Err(error) => return ExitCause::Failed(error),
};
let mut supervisor = Supervisor::new(Box::new(GuardTeardown::new(guard)));
let _subscription = match spawn_supervised_subscription(
&provider,
&config.underlying,
expiration_utc,
instruments,
&senders,
&mut supervisor,
)
.await
{
Ok(subscription) => subscription,
Err(error) => {
supervisor.fail(ChainViewError::provider(provider.id(), error));
return supervisor.run().await;
}
};
let (tx_events, mut rx_events) = event_channel();
let tick_child = supervisor.child_token();
let tick = spawn_tick_task(config.tick_interval, tx_events.clone(), tick_child.clone());
supervisor.register_ancillary(tick_child, Box::new(TokioTask::new(tick)));
let input_child = supervisor.child_token();
let input = spawn_input_reader(tx_events, input_child.clone());
supervisor.register_ancillary(input_child, Box::new(TokioTask::new(input)));
let backend = CrosstermBackend::new(io::stdout());
let terminal = match Terminal::new(backend) {
Ok(terminal) => terminal,
Err(error) => {
supervisor.fail(ChainViewError::Terminal(error.to_string()));
return supervisor.run().await;
}
};
let render_child = supervisor.child_token();
let root = supervisor.root_token();
let render = tokio::task::spawn_blocking(move || {
render_thread(terminal, &mut app, &mut bridge, &mut rx_events);
root.cancel();
});
supervisor.set_render(render_child, Box::new(TokioTask::new(render)));
drop(senders);
supervisor.run().await
}
fn render_thread(
mut terminal: Terminal<CrosstermBackend<Stdout>>,
app: &mut App,
bridge: &mut EventBridge,
rx_events: &mut tokio::sync::mpsc::Receiver<chainview::AppEvent>,
) {
let mut view = chainview::ViewState::new();
let _ = run_render_loop(
&mut terminal,
app,
bridge,
&mut view,
rx_events,
|_command| {},
);
}
fn run_replay(dir: PathBuf, config: Config) -> Result<(), ChainViewError> {
let runtime = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.map_err(|e| ChainViewError::Terminal(format!("tokio runtime: {e}")))?;
let cause = runtime.block_on(compose_and_run_replay(dir, config));
match cause {
ExitCause::Clean => Ok(()),
ExitCause::TaskPanicked => Err(ChainViewError::Terminal(
"a supervised task panicked; see the log".to_owned(),
)),
ExitCause::Failed(error) => Err(error),
}
}
async fn compose_and_run_replay(dir: PathBuf, config: Config) -> ExitCause {
let (mut bridge, senders) = EventBridge::new(config.channel_capacity);
let replay = ReplayState::new(dir.clone());
let mut app = App::new(
Mode::Replay(replay),
config.theme,
senders.tx_command.clone(),
)
.with_no_color(config.no_color);
let guard = match TerminalGuard::new() {
Ok(guard) => guard,
Err(error) => return ExitCause::Failed(error),
};
let mut supervisor = Supervisor::new(Box::new(GuardTeardown::new(guard)));
let (tx_events, mut rx_events) = event_channel();
let ceilings = ResourceCeilings::default();
let load_child = supervisor.child_token();
let load = spawn_bundle_load(dir, ceilings, tx_events.clone(), load_child.clone());
supervisor.register_ancillary(load_child, Box::new(TokioTask::new(load)));
let tick_child = supervisor.child_token();
let tick = spawn_tick_task(config.tick_interval, tx_events.clone(), tick_child.clone());
supervisor.register_ancillary(tick_child, Box::new(TokioTask::new(tick)));
let reload_tx = tx_events.clone();
let input_child = supervisor.child_token();
let input = spawn_input_reader(tx_events, input_child.clone());
supervisor.register_ancillary(input_child, Box::new(TokioTask::new(input)));
let backend = CrosstermBackend::new(io::stdout());
let terminal = match Terminal::new(backend) {
Ok(terminal) => terminal,
Err(error) => {
supervisor.fail(ChainViewError::Terminal(error.to_string()));
return supervisor.run().await;
}
};
let render_child = supervisor.child_token();
let root = supervisor.root_token();
let reload_cancel = supervisor.root_token();
let render = tokio::task::spawn_blocking(move || {
replay_render_thread(
terminal,
&mut app,
&mut bridge,
&mut rx_events,
ceilings,
reload_tx,
reload_cancel,
);
root.cancel();
});
supervisor.set_render(render_child, Box::new(TokioTask::new(render)));
drop(senders);
supervisor.run().await
}
fn replay_render_thread(
mut terminal: Terminal<CrosstermBackend<Stdout>>,
app: &mut App,
bridge: &mut EventBridge,
rx_events: &mut tokio::sync::mpsc::Receiver<chainview::AppEvent>,
ceilings: ResourceCeilings,
tx_events: tokio::sync::mpsc::Sender<chainview::AppEvent>,
cancel: tokio_util::sync::CancellationToken,
) {
let mut view = chainview::ViewState::new();
let _ = run_render_loop(
&mut terminal,
app,
bridge,
&mut view,
rx_events,
|command| route_replay_command(command, ceilings, &tx_events, &cancel),
);
}
fn route_replay_command(
command: DataCommand,
ceilings: ResourceCeilings,
tx_events: &tokio::sync::mpsc::Sender<chainview::AppEvent>,
cancel: &tokio_util::sync::CancellationToken,
) {
match command {
DataCommand::ReloadBundle(dir) => {
let _load = spawn_bundle_load(dir, ceilings, tx_events.clone(), cancel.child_token());
}
DataCommand::Subscribe { .. }
| DataCommand::Unsubscribe { .. }
| DataCommand::Reconnect
| DataCommand::Rediscover => {}
}
}
fn positive_or_one(value: f64) -> Positive {
Positive::new(value).unwrap_or_else(|_| positive_one())
}
fn positive_one() -> Positive {
Positive::new(1.0).unwrap_or(Positive::ZERO)
}
fn empty_seed(underlying: &str, source: &SourceBinding, now: DateTime<Utc>) -> ChainFetch {
let chain = OptionChain::new(underlying, positive_one(), now.to_rfc3339(), None, None);
ChainFetch::new(
chain,
ExpirySource::new(underlying, now, source.provider.clone()),
AliasCatalog::new(),
)
}
fn now_utc() -> DateTime<Utc> {
let since = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or(Duration::ZERO);
let secs = i64::try_from(since.as_secs()).unwrap_or(i64::MAX);
DateTime::<Utc>::from_timestamp(secs, since.subsec_nanos()).unwrap_or(DateTime::<Utc>::MIN_UTC)
}
fn validate_replay_dir(dir: &Path) -> Result<(), ChainViewError> {
match std::fs::metadata(dir) {
Ok(meta) if meta.is_dir() => Ok(()),
Ok(_) => Err(
BundleError::Io(format!("replay path is not a directory: {}", dir.display())).into(),
),
Err(_) => Err(BundleError::Io(format!(
"replay bundle directory not found: {}",
dir.display()
))
.into()),
}
}
#[cfg(test)]
mod tests {
use super::{Cli, route_replay_command, validate_replay_dir};
use chainview::config::ModeSelect;
use chainview::{
AppEvent, BundleLoadResult, ChainViewError, Command as DataCommand, ResourceCeilings,
};
use clap::Parser;
use std::path::PathBuf;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
fn mode_of(args: &[&str]) -> Result<ModeSelect, clap::Error> {
Cli::try_parse_from(args).map(|cli| cli.into_overrides().mode)
}
#[test]
fn test_cli_no_subcommand_selects_live() {
match mode_of(&["chainview"]) {
Ok(mode) => assert_eq!(mode, ModeSelect::Live),
Err(e) => panic!("expected Live, got parse error: {e}"),
}
}
#[test]
fn test_cli_replay_subcommand_selects_replay_with_dir() {
match mode_of(&["chainview", "replay", "./run-2026-07-01/"]) {
Ok(mode) => assert_eq!(mode, ModeSelect::Replay(PathBuf::from("./run-2026-07-01/"))),
Err(e) => panic!("expected Replay, got parse error: {e}"),
}
}
#[test]
fn test_cli_replay_ignores_live_only_flags() {
match mode_of(&[
"chainview",
"--provider",
"ig",
"--underlying",
"SPY",
"replay",
"./bundle/",
]) {
Ok(mode) => assert_eq!(mode, ModeSelect::Replay(PathBuf::from("./bundle/"))),
Err(e) => panic!("expected Replay, got parse error: {e}"),
}
}
#[test]
fn test_cli_replay_requires_a_directory() {
assert!(mode_of(&["chainview", "replay"]).is_err());
}
#[test]
fn test_cli_replay_rejects_extra_positional() {
assert!(mode_of(&["chainview", "replay", "./a", "./b"]).is_err());
}
#[test]
fn test_validate_replay_dir_accepts_an_existing_directory() {
let dir = match tempfile::tempdir() {
Ok(d) => d,
Err(e) => panic!("failed to make a temp dir: {e}"),
};
assert!(validate_replay_dir(dir.path()).is_ok());
}
#[test]
fn test_validate_replay_dir_rejects_a_missing_directory() {
let dir = match tempfile::tempdir() {
Ok(d) => d,
Err(e) => panic!("failed to make a temp dir: {e}"),
};
let missing = dir.path().join("does-not-exist");
match validate_replay_dir(&missing) {
Err(ChainViewError::Bundle(_)) => {}
other => panic!("expected a friendly Bundle error, got {other:?}"),
}
}
#[test]
fn test_validate_replay_dir_rejects_a_file() {
use std::io::Write;
let dir = match tempfile::tempdir() {
Ok(d) => d,
Err(e) => panic!("failed to make a temp dir: {e}"),
};
let file_path = dir.path().join("manifest.json");
match std::fs::File::create(&file_path).and_then(|mut f| f.write_all(b"{}")) {
Ok(()) => {}
Err(e) => panic!("failed to write a temp file: {e}"),
}
match validate_replay_dir(&file_path) {
Err(ChainViewError::Bundle(_)) => {}
other => panic!("expected a friendly Bundle error for a file, got {other:?}"),
}
}
#[tokio::test]
async fn test_route_replay_command_reload_respawns_load_and_emits_bundle_loaded() {
let (tx, mut rx) = mpsc::channel::<AppEvent>(8);
let cancel = CancellationToken::new();
let dir = std::env::temp_dir().join("chainview-route-missing-bundle-34");
let _ = std::fs::remove_dir_all(&dir);
route_replay_command(
DataCommand::ReloadBundle(dir),
ResourceCeilings::default(),
&tx,
&cancel,
);
match tokio::time::timeout(std::time::Duration::from_secs(5), rx.recv()).await {
Ok(Some(AppEvent::BundleLoaded(BundleLoadResult::Failed(message)))) => {
assert!(
!message.is_empty(),
"the failure carries a non-secret message"
);
}
other => {
panic!("expected BundleLoaded(Failed) from the respawned load, got {other:?}")
}
}
}
#[tokio::test]
async fn test_route_replay_command_ignores_live_only_commands() {
let (tx, mut rx) = mpsc::channel::<AppEvent>(8);
let cancel = CancellationToken::new();
for command in [DataCommand::Reconnect, DataCommand::Rediscover] {
route_replay_command(command, ResourceCeilings::default(), &tx, &cancel);
}
assert!(
rx.try_recv().is_err(),
"live-only commands produce no replay event and spawn no load",
);
}
}