use std::path::PathBuf; use std::sync::Arc; use std::time::Duration; use anyhow::Context; use clap::Parser; use tokio::process::Command; mod cli; mod config; mod gateway; mod imap; mod jmap; mod maildir; mod mailparse; mod smtp; mod store; mod worker; use store::{JobRecord, Store, UserRecord}; #[tokio::main] async fn main() -> anyhow::Result<()> { let args = cli::Args::parse(); // Tracing first, so Config::load's "config not found" warning lands in the // log instead of being emitted into the void before the subscriber exists. tracing_subscriber::fmt() .with_env_filter( tracing_subscriber::EnvFilter::try_from_default_env() .unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")), ) .init(); let cfg = config::Config::load(&args.config) .await .context("failed to load configuration")?; match args.command.unwrap_or_default() { cli::Command::Worker { job_id } => run_worker(&cfg, job_id).await, cli::Command::Serve => run_server(&args.config, &cfg).await, } } /// Child-process entry: process one LLM job and exit. The job is already /// claimed (status=running) by the parent dispatcher; we just run it. async fn run_worker(cfg: &config::Config, job_id: i32) -> anyhow::Result<()> { let data_dir = PathBuf::from(&cfg.daemon.data_dir); let users: Vec = cfg .users .iter() .map(|u| UserRecord { email: u.email.clone(), password: u.password.clone(), }) .collect(); let provider_names: Vec = cfg.llm.providers.iter().map(|p| p.name.clone()).collect(); let store = Store::new(data_dir, provider_names, users) .context("worker: failed to open store")?; let worker = worker::Worker::new(Arc::new(store), cfg.llm.clone(), cfg.daemon.domain.clone()); worker.run_single_job(job_id).await } /// Server entry: SMTP + IMAP plus a dispatcher that spawns one child worker /// process per claimed LLM job. async fn run_server(config_path: &std::path::Path, cfg: &config::Config) -> anyhow::Result<()> { let data_dir = PathBuf::from(&cfg.daemon.data_dir); std::fs::create_dir_all(&data_dir) .with_context(|| format!("create data dir {}", data_dir.display()))?; let users: Vec = cfg .users .iter() .map(|u| UserRecord { email: u.email.clone(), password: u.password.clone(), }) .collect(); let provider_names: Vec = cfg.llm.providers.iter().map(|p| p.name.clone()).collect(); let store = Store::new(data_dir, provider_names, users) .context("failed to open store")?; let store = Arc::new(store); let smtp_server = smtp::SmtpServer::new( cfg.smtp.clone(), cfg.daemon.clone(), store.clone(), ); smtp_server .listen() .await .context("SMTP server failed to start")?; let imap_server = imap::ImapServer::new(cfg.imap.clone(), store.clone()); imap_server .listen() .await .context("IMAP server failed to start")?; let jmap_server = jmap::JmapServer::new(cfg.jmap.clone(), store.clone()); jmap_server .listen() .await .context("JMAP server failed to start")?; let dispatcher = tokio::spawn(dispatch_loop(store.clone(), config_path.to_path_buf())); tracing::info!( domain = %cfg.daemon.domain, providers = ?cfg.llm.providers.iter().map(|p| p.name.as_str()).collect::>(), "posthorn ready" ); tokio::select! { _ = tokio::signal::ctrl_c() => {} _ = dispatcher => { tracing::warn!("job dispatcher exited unexpectedly"); } } tracing::info!("posthorn stopping"); Ok(()) } /// Poll the job queue; for each newly-claimed job, spawn a child /// `posthorn worker --job N` process. Children run concurrently and are /// detached — we log their exit but never block the poll loop on them. A /// bounded semaphore caps in-flight children so a backlog can't fork-bomb. async fn dispatch_loop(store: Arc, config_path: PathBuf) { // Cap concurrent worker children. Each drives a generation on the GPU // gateway, so this is more about not over-subscribing the model than RAM. const MAX_CONCURRENT: usize = 4; let semaphore = Arc::new(tokio::sync::Semaphore::new(MAX_CONCURRENT)); let exe = match std::env::current_exe() { Ok(e) => e, Err(e) => { tracing::error!(error = %e, "dispatcher: cannot resolve current_exe; exiting"); return; } }; // Generous: the app's whole point is patience — but not forever. A hung // gateway (e.g. a provider that accepts the connection and never responds) // would otherwise hold a semaphore permit indefinitely. const JOB_DEADLINE: Duration = Duration::from_secs(60 * 60); // Reap window MUST exceed the deadline: orphaned children of a crashed // daemon keep running to completion, and reaping too eagerly would // requeue a job whose orphan later finishes -> duplicate reply. const STALE: Duration = Duration::from_secs(90 * 60); // Requeue anything wedged `running` by a previous crash before we start. if let Err(e) = store.requeue_stale_jobs(STALE).await { tracing::warn!(error = %e, "initial requeue of stale jobs failed"); } let mut interval = tokio::time::interval(Duration::from_secs(5)); loop { interval.tick().await; if let Err(e) = store.requeue_stale_jobs(STALE).await { tracing::warn!(error = %e, "requeue of stale jobs failed"); } let jobs: Vec = match store.claim_pending_jobs(4).await { Ok(jobs) => jobs, Err(e) => { tracing::warn!(error = %e, "claim pending jobs failed"); continue; } }; for job in jobs { // Acquire a permit before forking; if we're saturated, hold the // job (already claimed as running) until a slot frees. let permit = semaphore .clone() .acquire_owned() .await .expect("semaphore closed"); let exe = exe.clone(); let config_path = config_path.clone(); let store = store.clone(); let job_id = job.id; let thread_id = job.thread_id; let message_id = job.message_id.clone(); tokio::spawn(async move { tracing::info!(job_id, thread_id, %message_id, "spawning worker child"); let mut child = Command::new(&exe) .arg("--config") .arg(&config_path) .arg("worker") .arg("--job-id") .arg(job_id.to_string()) .stdin(std::process::Stdio::null()) .stdout(std::process::Stdio::inherit()) .stderr(std::process::Stdio::inherit()) .spawn(); match child.as_mut() { Ok(child) => match tokio::time::timeout(JOB_DEADLINE, child.wait()).await { Ok(Ok(status)) if status.success() => { tracing::info!(job_id, "worker child exited cleanly"); } Ok(Ok(status)) => { tracing::warn!(job_id, %status, "worker child exited non-zero"); let err = format!("worker exited {status}"); let _ = store.fail_job(job_id, &err).await; } Ok(Err(e)) => { tracing::warn!(job_id, error = %e, "worker child wait failed"); let err = format!("wait failed: {e}"); let _ = store.fail_job(job_id, &err).await; } Err(_) => { tracing::warn!(job_id, "worker exceeded deadline; killing"); let _ = child.kill().await; let _ = store .fail_job(job_id, "generation exceeded deadline; killed") .await; } }, Err(e) => { tracing::error!(job_id, error = %e, "failed to spawn worker child"); let err = format!("spawn failed: {e}"); let _ = store.fail_job(job_id, &err).await; } } drop(permit); }); } } }