mod app;
mod broker;
mod config;
mod error;
mod models;
mod ui;
mod update;
mod utils;
use anyhow::Result;
use clap::Parser;
use crossterm::{
event::{DisableMouseCapture, EnableMouseCapture},
execute,
terminal::{disable_raw_mode, enable_raw_mode, EnterAlternateScreen, LeaveAlternateScreen},
};
use ratatui::{backend::CrosstermBackend, Terminal};
use std::{io, time::Duration};
use tokio::time;
use crate::app::App;
use crate::broker::{create_broker, Broker};
use crate::config::Config;
use crate::ui::events::{handle_key_event, next_event, AppEvent};
use clap::Subcommand;
#[derive(Parser, Debug)]
#[command(author, version, about, long_about = None)]
struct Cli {
#[command(subcommand)]
command: Option<Commands>,
#[arg(short, long, global = true)]
broker: Option<String>,
#[arg(long, global = true)]
result_backend: Option<String>,
#[arg(short, long, global = true)]
config: Option<std::path::PathBuf>,
}
#[derive(Subcommand, Debug)]
enum Commands {
Init,
Config,
SetBroker {
url: String,
},
SetRefresh {
interval: u64,
},
}
#[tokio::main]
async fn main() -> Result<()> {
let cli = Cli::parse();
match cli.command {
Some(Commands::Init) => {
run_init_command().await?;
return Ok(());
}
Some(Commands::Config) => {
show_config()?;
return Ok(());
}
Some(Commands::SetBroker { url }) => {
set_broker_url(&url)?;
return Ok(());
}
Some(Commands::SetRefresh { interval }) => {
set_refresh_interval(interval)?;
return Ok(());
}
None => {
run_tui_app(cli.broker, cli.config).await?;
}
}
Ok(())
}
async fn run_tui_app(
broker_arg: Option<String>,
config_arg: Option<std::path::PathBuf>,
) -> Result<()> {
let config = if let Some(config_path) = config_arg {
Config::from_file(config_path)?
} else {
Config::load_or_create_default()?
};
let current_version = env!("CARGO_PKG_VERSION");
tokio::spawn(async move {
if let Some(update) = update::check_for_update(current_version).await {
update.print_notification();
}
});
let broker_url = broker_arg.unwrap_or_else(|| config.broker.url.clone());
let broker: Box<dyn Broker> = match create_broker(&broker_url).await {
Ok(broker) => broker,
Err(e) => {
let (broker_type, url_hint) = if broker_url.starts_with("redis://") {
("Redis", "redis://localhost:6379/0")
} else if broker_url.starts_with("amqp://") {
("RabbitMQ", "amqp://guest:guest@localhost:5672//")
} else {
(
"Unknown",
"redis://localhost:6379/0 or amqp://localhost:5672//",
)
};
eprintln!("\n❌ Failed to connect to {broker_type} broker at {broker_url}");
eprintln!("\n{e}");
eprintln!("\n📋 Quick Setup Guide:");
eprintln!("1. For Redis:");
eprintln!(" - Docker: docker run -d -p 6379:6379 redis");
eprintln!(" - macOS: brew services start redis");
eprintln!(" - Verify: redis-cli ping");
eprintln!("\n2. For RabbitMQ:");
eprintln!(" - Docker: docker run -d -p 5672:5672 rabbitmq");
eprintln!(" - Verify: amqp://guest:guest@localhost:5672//");
eprintln!("\n3. Run lazycelery:");
eprintln!(" lazycelery --broker {url_hint}");
eprintln!("\n💡 For more help: https://github.com/Fgudes90/lazycelery");
std::process::exit(1);
}
};
let mut app = App::new(broker);
enable_raw_mode()?;
let mut stdout = io::stdout();
execute!(stdout, EnterAlternateScreen, EnableMouseCapture)?;
let backend = CrosstermBackend::new(stdout);
let mut terminal = Terminal::new(backend)?;
let res = run_app(&mut terminal, &mut app, &config).await;
disable_raw_mode()?;
execute!(
terminal.backend_mut(),
LeaveAlternateScreen,
DisableMouseCapture
)?;
terminal.show_cursor()?;
if let Err(err) = res {
eprintln!("Error: {err}");
}
Ok(())
}
async fn run_app(
terminal: &mut Terminal<CrosstermBackend<io::Stdout>>,
app: &mut App,
config: &Config,
) -> Result<()> {
app.refresh_data().await?;
let mut refresh_interval = time::interval(Duration::from_millis(config.ui.refresh_interval));
let tick_rate = Duration::from_millis(50);
loop {
terminal.draw(|f| ui::draw(f, app))?;
tokio::select! {
event = next_event(tick_rate) => {
match event? {
AppEvent::Key(key) => {
let should_execute = app.show_confirmation && matches!(
key.code,
crossterm::event::KeyCode::Char('y') |
crossterm::event::KeyCode::Char('Y') |
crossterm::event::KeyCode::Enter
);
handle_key_event(key, app);
if should_execute {
app.execute_pending_action().await?;
}
if app.should_quit {
return Ok(());
}
}
AppEvent::Tick => {}
AppEvent::Refresh => {
app.refresh_data().await?;
}
}
}
_ = refresh_interval.tick() => {
app.refresh_data().await?;
}
}
}
}
async fn run_init_command() -> Result<()> {
use std::io::{self, Write};
println!("🚀 Welcome to LazyCelery Setup!\n");
let config_dir = dirs::config_dir()
.ok_or_else(|| anyhow::anyhow!("Could not determine config directory"))?
.join("lazycelery");
let config_path = config_dir.join("config.toml");
if config_path.exists() {
print!(
"⚠️ Configuration already exists at {}. Overwrite? (y/N): ",
config_path.display()
);
io::stdout().flush()?;
let mut input = String::new();
io::stdin().read_line(&mut input)?;
if !input.trim().eq_ignore_ascii_case("y") {
println!("❌ Setup cancelled.");
return Ok(());
}
}
print!("📡 Enter your Celery broker URL (default: redis://localhost:6379/0): ");
io::stdout().flush()?;
let mut broker_url = String::new();
io::stdin().read_line(&mut broker_url)?;
let broker_url = broker_url.trim();
let broker_url = if broker_url.is_empty() {
"redis://localhost:6379/0"
} else {
broker_url
};
if !broker_url.starts_with("redis://") && !broker_url.starts_with("amqp://") {
eprintln!("❌ Invalid broker URL. Must start with redis:// or amqp://");
return Ok(());
}
print!("🔄 Enter UI refresh interval in milliseconds (default: 1000): ");
io::stdout().flush()?;
let mut refresh_input = String::new();
io::stdin().read_line(&mut refresh_input)?;
let refresh_interval: u64 = refresh_input.trim().parse().unwrap_or(1000);
let config = Config {
broker: crate::config::BrokerConfig {
url: broker_url.to_string(),
timeout: 30,
retry_attempts: 3,
},
ui: crate::config::UiConfig {
refresh_interval,
theme: "dark".to_string(),
},
};
std::fs::create_dir_all(&config_dir)?;
let toml_string = toml::to_string_pretty(&config)?;
std::fs::write(&config_path, toml_string)?;
println!("\n✅ Configuration saved to: {}", config_path.display());
println!("\n📋 You can now run 'lazycelery' to start monitoring your Celery workers!");
print!("\n🔌 Test connection to broker? (Y/n): ");
io::stdout().flush()?;
let mut test_input = String::new();
io::stdin().read_line(&mut test_input)?;
if !test_input.trim().eq_ignore_ascii_case("n") {
print!("🔄 Testing connection to {}... ", config.broker.url);
io::stdout().flush()?;
match test_broker_connection(&config.broker.url).await {
Ok(_) => println!("✅ Success!"),
Err(e) => println!("❌ Failed: {e}"),
}
}
Ok(())
}
fn show_config() -> Result<()> {
let config_dir = dirs::config_dir()
.ok_or_else(|| anyhow::anyhow!("Could not determine config directory"))?
.join("lazycelery");
let config_path = config_dir.join("config.toml");
if !config_path.exists() {
eprintln!("❌ No configuration found. Run 'lazycelery init' to create one.");
return Ok(());
}
let config = Config::from_file(config_path.clone())?;
println!("📋 Current Configuration");
println!("📍 Location: {}", config_path.display());
println!("\n[broker]");
println!(" url = \"{}\"", config.broker.url);
println!(" timeout = {}", config.broker.timeout);
println!(" retry_attempts = {}", config.broker.retry_attempts);
println!("\n[ui]");
println!(" refresh_interval = {}", config.ui.refresh_interval);
println!(" theme = \"{}\"", config.ui.theme);
Ok(())
}
fn set_broker_url(url: &str) -> Result<()> {
if !url.starts_with("redis://") && !url.starts_with("amqp://") {
eprintln!("❌ Invalid broker URL. Must start with redis:// or amqp://");
return Ok(());
}
let config_dir = dirs::config_dir()
.ok_or_else(|| anyhow::anyhow!("Could not determine config directory"))?
.join("lazycelery");
let config_path = config_dir.join("config.toml");
let mut config = if config_path.exists() {
Config::from_file(config_path.clone())?
} else {
std::fs::create_dir_all(&config_dir)?;
Config::default()
};
config.broker.url = url.to_string();
let toml_string = toml::to_string_pretty(&config)?;
std::fs::write(&config_path, toml_string)?;
println!("✅ Broker URL updated to: {url}");
println!("📍 Configuration saved to: {}", config_path.display());
Ok(())
}
fn set_refresh_interval(interval: u64) -> Result<()> {
let config_dir = dirs::config_dir()
.ok_or_else(|| anyhow::anyhow!("Could not determine config directory"))?
.join("lazycelery");
let config_path = config_dir.join("config.toml");
let mut config = if config_path.exists() {
Config::from_file(config_path.clone())?
} else {
std::fs::create_dir_all(&config_dir)?;
Config::default()
};
config.ui.refresh_interval = interval;
let toml_string = toml::to_string_pretty(&config)?;
std::fs::write(&config_path, toml_string)?;
println!("✅ Refresh interval updated to: {interval}ms");
println!("📍 Configuration saved to: {}", config_path.display());
Ok(())
}
async fn test_broker_connection(url: &str) -> Result<()> {
create_broker(url).await?;
Ok(())
}