From 7319c8715e02aa7fba4619b8b56a6d502d241396 Mon Sep 17 00:00:00 2001 From: Luke Street Date: Fri, 2 Jan 2026 19:27:48 -0700 Subject: [PATCH] Integrate apalis-board --- Cargo.lock | 111 +++++++++++++++++++++++++++++++++----- Cargo.toml | 8 +-- crates/core/src/config.rs | 3 +- crates/jobs/Cargo.toml | 3 +- crates/jobs/src/lib.rs | 19 +++++-- crates/web/Cargo.toml | 1 + crates/web/src/main.rs | 100 +++++++++++++++++++++++++--------- 7 files changed, 199 insertions(+), 46 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 8af0142..f177dc8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -93,9 +93,9 @@ checksum = "a23eb6b1614318a8071c9b2521f36b424b2c83db5eb3a0fead4a6c0809af6e61" [[package]] name = "apalis" -version = "1.0.0-beta.2" +version = "1.0.0-rc.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c6556b89bb5cb40dab6e4d7ee1f2abfef6d372bd3aef300a3faa992bcb13bda8" +checksum = "f93be0eb33b912f5e66004d0b756423c285273259068b1c80a71d7842658189b" dependencies = [ "apalis-core", "futures-util", @@ -106,10 +106,59 @@ dependencies = [ ] [[package]] -name = "apalis-core" -version = "1.0.0-beta.2" +name = "apalis-board" +version = "1.0.0-rc.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "49039d4eb476cc05e196210153dbb34be180de9a0ef8b7a666fde0f80c5c357d" +checksum = "eae7f45f6f1150f5deac364df7a118e3087adbc18a78267b5b93e143cbf2ab0d" +dependencies = [ + "apalis-board-api", +] + +[[package]] +name = "apalis-board-api" +version = "1.0.0-rc.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "73630d08d0c01b1b9f4d74f22171b53a904289f1907e606d41759566f6ac642f" +dependencies = [ + "apalis-board-types", + "apalis-core", + "axum", + "futures", + "include_dir", + "serde", + "serde_json", + "thiserror 2.0.17", + "tokio", + "tracing-core", + "tracing-subscriber", +] + +[[package]] +name = "apalis-board-types" +version = "1.0.0-rc.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a7cd2f7fb89d1aeb2e2f1c4a9c35e23a394566e957bd634074c2180a55a97c1c" +dependencies = [ + "serde", + "thiserror 2.0.17", +] + +[[package]] +name = "apalis-codec" +version = "0.1.0-rc.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a5ed6bb8e64c360ed4ad666a6cbc42e9e6df73087461dc4071f510a3af284637" +dependencies = [ + "apalis-core", + "serde", + "serde_json", +] + +[[package]] +name = "apalis-core" +version = "1.0.0-rc.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6b1557d680ee4a9b42a76ab3a9572cee1a00d45e7eb455427d906c42774766e7" dependencies = [ "futures-channel", "futures-core", @@ -118,7 +167,6 @@ dependencies = [ "futures-util", "pin-project", "serde", - "serde_json", "thiserror 2.0.17", "tower-layer", "tower-service", @@ -127,9 +175,9 @@ dependencies = [ [[package]] name = "apalis-cron" -version = "1.0.0-beta.2" +version = "1.0.0-rc.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "53230aa8a09cdbe2e254d1b59999efe9f1cad09b66d48475f51b3fb38be6b840" +checksum = "d82cc92dc7b6c48c29f308379c6095d82738ebe0831d9bfb27b8a1a8988ee981" dependencies = [ "apalis-core", "chrono", @@ -141,9 +189,9 @@ dependencies = [ [[package]] name = "apalis-sql" -version = "1.0.0-beta.2" +version = "1.0.0-rc.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "604ed96ae8ff20d4c6f8533ffc9f8c91c5f99da6bc564de6a37b94829522f9ec" +checksum = "5ade5d8faa60e9975b01d3bb1ebc5028589aa4986365eaa4d080d30ed3b5141f" dependencies = [ "apalis-core", "chrono", @@ -154,10 +202,11 @@ dependencies = [ [[package]] name = "apalis-sqlite" -version = "1.0.0-beta.4" +version = "1.0.0-rc.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8be443bddc6ba6e023b8364304def0a6e3428f258b5addf5aa74bcb5688eb3d1" +checksum = "fd43020ce13d6cb8c8d8c09657a6691d8490cd1ce0d8bc0f7fef8bf9b23cfe86" dependencies = [ + "apalis-codec", "apalis-core", "apalis-sql", "chrono", @@ -174,9 +223,9 @@ dependencies = [ [[package]] name = "apalis-workflow" -version = "0.1.0-beta.2" +version = "0.1.0-rc.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "afa1869ce395c0b727315a2be39508eb77ef94ed75492e7b43334b13241b1afa" +checksum = "bc024da2d5d3ab59cc9fea099a2e2b20de5ff608f2e287abcb73aa45e4966a89" dependencies = [ "apalis-core", "futures", @@ -1104,6 +1153,7 @@ version = "0.1.0" dependencies = [ "anyhow", "apalis", + "apalis-codec", "apalis-cron", "apalis-sqlite", "apalis-workflow", @@ -1125,6 +1175,7 @@ version = "0.1.0" dependencies = [ "anyhow", "apalis", + "apalis-board", "axum", "axum_typed_multipart", "bytes", @@ -2160,6 +2211,25 @@ version = "1.12.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e7c5cedc30da3a610cac6b4ba17597bdf7152cf974e8aab3afb3d54455e371c8" +[[package]] +name = "include_dir" +version = "0.7.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "923d117408f1e49d914f1a379a309cffe4f18c05cf4e3d12e613a15fc81bd0dd" +dependencies = [ + "include_dir_macros", +] + +[[package]] +name = "include_dir_macros" +version = "0.7.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7cab85a7ed0bd5f0e76d93846e0147172bed2e2d3f859bcc33a8d9699cad1a75" +dependencies = [ + "proc-macro2", + "quote", +] + [[package]] name = "indexmap" version = "2.12.1" @@ -4925,6 +4995,16 @@ dependencies = [ "tracing-core", ] +[[package]] +name = "tracing-serde" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "704b1aeb7be0d0a84fc9828cae51dab5970fee5088f83d1dd7ee6f6246fc6ff1" +dependencies = [ + "serde", + "tracing-core", +] + [[package]] name = "tracing-subscriber" version = "0.3.22" @@ -4935,12 +5015,15 @@ dependencies = [ "nu-ansi-term", "once_cell", "regex-automata", + "serde", + "serde_json", "sharded-slab", "smallvec", "thread_local", "tracing", "tracing-core", "tracing-log", + "tracing-serde", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index 0aa54df..86c9c2f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -18,9 +18,11 @@ rust-version = "1.88" [workspace.dependencies] anyhow = "1.0" -apalis = "1.0.0-beta.2" -apalis-cron = "1.0.0-beta.2" -apalis-sqlite = "1.0.0-beta.4" +apalis = "1.0.0-rc.1" +apalis-board = { version = "1.0.0-rc.1", features = ["axum"] } +apalis-codec = "0.1.0-rc.1" +apalis-cron = "1.0.0-rc.1" +apalis-sqlite = "1.0.0-rc.1" axum = { version = "0.8", features = ["macros"] } futures-util = "0.3.31" hex = "0.4" diff --git a/crates/core/src/config.rs b/crates/core/src/config.rs index 7f48da2..b13f6b2 100644 --- a/crates/core/src/config.rs +++ b/crates/core/src/config.rs @@ -13,6 +13,7 @@ pub struct Config { #[derive(Debug, Clone, Deserialize, Serialize)] pub struct ServerConfig { pub port: u16, + pub jobs_port: Option, #[serde(default)] pub dev_mode: bool, } @@ -64,6 +65,6 @@ pub struct WorkerConfig { impl Default for WorkerConfig { fn default() -> Self { - Self { workflow_run_concurrency: 3, refresh_project_concurrency: 1, retry_attempts: 5 } + Self { workflow_run_concurrency: 3, refresh_project_concurrency: 3, retry_attempts: 5 } } } diff --git a/crates/jobs/Cargo.toml b/crates/jobs/Cargo.toml index 9e87dd1..aa846ee 100644 --- a/crates/jobs/Cargo.toml +++ b/crates/jobs/Cargo.toml @@ -6,9 +6,10 @@ publish = false [dependencies] anyhow.workspace = true -apalis.workspace = true +apalis-codec.workspace = true apalis-cron.workspace = true apalis-sqlite.workspace = true +apalis.workspace = true apalis-workflow = "0.1.0-beta.2" decomp-dev-core = { path = "../core" } decomp-dev-db = { path = "../db" } diff --git a/crates/jobs/src/lib.rs b/crates/jobs/src/lib.rs index c02e8a2..2706982 100644 --- a/crates/jobs/src/lib.rs +++ b/crates/jobs/src/lib.rs @@ -10,6 +10,7 @@ use apalis::{ }, prelude::*, }; +use apalis_codec::json::JsonCodec; use apalis_sqlite::{CompactType, SqliteStorage, fetcher::SqliteFetcher}; use decomp_dev_core::config::{Config, DbConfig, WorkerConfig}; use decomp_dev_db::Database; @@ -28,7 +29,7 @@ pub struct JobContext { } /// Type alias for the default codec used by SqliteStorage. -type DefaultCodec = json::JsonCodec; +type DefaultCodec = JsonCodec; /// Type alias for workflow run storage. pub type WorkflowRunStorage = SqliteStorage; @@ -55,8 +56,8 @@ impl JobStorage { SqlitePool::connect(&db.jobs_url).await.context("Failed to connect to database")?; SqliteStorage::setup(&pool).await?; Ok(Arc::new(Self { - workflow_run: SqliteStorage::new(&pool), - refresh_project: SqliteStorage::new(&pool), + workflow_run: create_storage(&pool), + refresh_project: create_storage(&pool), })) } @@ -67,6 +68,18 @@ impl JobStorage { pub fn refresh_project(&self) -> RefreshProjectStorage { self.refresh_project.clone() } } +fn create_storage(pool: &SqlitePool) -> SqliteStorage { + let config = apalis_sqlite::Config::new(std::any::type_name::()).with_poll_interval( + StrategyBuilder::new() + .apply( + IntervalStrategy::new(Duration::from_millis(100)) + .with_backoff(BackoffConfig::new(Duration::from_secs(1))), + ) + .build(), + ); + SqliteStorage::new_with_config(pool, &config) +} + /// Create the job monitor with all workers. pub fn create_monitor( storage: Arc, diff --git a/crates/web/Cargo.toml b/crates/web/Cargo.toml index 4c6f9fd..7034d45 100644 --- a/crates/web/Cargo.toml +++ b/crates/web/Cargo.toml @@ -8,6 +8,7 @@ default-run = "decomp-dev-web" [dependencies] anyhow.workspace = true apalis.workspace = true +apalis-board.workspace = true axum.workspace = true axum_typed_multipart = "0.16" decomp-dev-auth = { path = "../auth" } diff --git a/crates/web/src/main.rs b/crates/web/src/main.rs index 31a4984..fa32dd4 100644 --- a/crates/web/src/main.rs +++ b/crates/web/src/main.rs @@ -9,12 +9,18 @@ use std::{ io::BufReader, net::{IpAddr, Ipv4Addr, SocketAddr}, str::FromStr, - sync::Arc, + sync::{Arc, Mutex}, time::Duration, }; +use anyhow::Context; +use apalis_board::axum::{ + framework::{ApiBuilder, RegisterRoute}, + sse::{TracingBroadcaster, TracingSubscriber}, + ui::ServeUI, +}; use axum::{ - Router, + Extension, Router, extract::{ConnectInfo, FromRef}, http::{Method, Request, StatusCode, header}, middleware, @@ -26,15 +32,18 @@ use decomp_dev_jobs::{JobContext, JobStorage, create_monitor}; use tokio::{net::TcpListener, signal}; use tower::ServiceBuilder; use tower_http::{ - ServiceBuilderExt, cors, - cors::CorsLayer, + ServiceBuilderExt, + cors::{self, CorsLayer}, + normalize_path::NormalizePathLayer, timeout::TimeoutLayer, trace::{DefaultOnResponse, MakeSpan, TraceLayer}, }; use tower_sessions::{Expiry, SessionManagerLayer, SessionStore, cookie::SameSite}; use tower_sessions_sqlx_store::SqliteStore; use tracing::{Level, Span}; -use tracing_subscriber::{EnvFilter, filter::LevelFilter}; +use tracing_subscriber::{ + EnvFilter, Layer, filter::LevelFilter, layer::SubscriberExt, util::SubscriberInitExt, +}; use crate::handlers::{build_router, csp::csp_middleware}; @@ -48,13 +57,14 @@ pub struct AppState { #[tokio::main] async fn main() { - tracing_subscriber::fmt() - .with_env_filter( - EnvFilter::builder() - // Default to info level - .with_default_directive(LevelFilter::INFO.into()) - .from_env_lossy(), - ) + let broadcaster = TracingBroadcaster::create(); + let env_filter = EnvFilter::builder() + // Default to info level + .with_default_directive(LevelFilter::INFO.into()) + .from_env_lossy(); + tracing_subscriber::registry() + .with(tracing_subscriber::fmt::layer().with_filter(env_filter.clone())) + .with(TracingSubscriber::new(&broadcaster).layer()) .init(); let config: Arc = { @@ -86,6 +96,8 @@ async fn main() { // Build the router let port = state.config.server.port; + let jobs_port = state.config.server.jobs_port; + let jobs_router = jobs_app(&state.jobs, broadcaster); let router = app(state, session_store).into_make_service_with_connect_info::(); // Create the listener @@ -97,7 +109,7 @@ async fn main() { let fds = libsystemd::activation::receive_descriptors_with_names(false) .expect("Failed to receive fds"); if let Some((fd, name)) = fds.into_iter().next() { - tracing::info!("Listening on {}", name); + tracing::info!("Web server: Listening on {}", name); let std_listener = unsafe { std::net::TcpListener::from_raw_fd(fd.into_raw_fd()) }; std_listener.set_nonblocking(true).expect("Failed to set non-blocking"); listener = @@ -108,10 +120,21 @@ async fn main() { Some(listener) => listener, None => { let addr = SocketAddr::from((Ipv4Addr::UNSPECIFIED, port)); - tracing::info!("Listening on {}", addr); + tracing::info!("Web server: Listening on {}", addr); TcpListener::bind(addr).await.expect("bind error") } }; + let jobs_listener = match jobs_port { + Some(port) => { + let addr = SocketAddr::from((Ipv4Addr::LOCALHOST, port)); + tracing::info!("Jobs server: Listening on {}", addr); + Some(TcpListener::bind(addr).await.expect("bind error")) + } + None => { + tracing::info!("Jobs server: Disabled"); + None + } + }; #[cfg(target_os = "linux")] { @@ -121,20 +144,34 @@ async fn main() { // Run both the web server and job monitor concurrently, with graceful shutdown let web_server = async { - axum::serve(listener, router) + let result = axum::serve(listener, router) .with_graceful_shutdown(shutdown_signal()) .await - .map_err(|e| anyhow::anyhow!("Web server error: {e}")) + .context("Web server error"); + tracing::info!("Web server stopped"); + result + }; + let jobs_server = async { + if let Some(jobs_listener) = jobs_listener { + let result = axum::serve(jobs_listener, jobs_router) + .with_graceful_shutdown(shutdown_signal()) + .await + .context("Jobs server error"); + tracing::info!("Jobs server stopped"); + result + } else { + Ok(()) + } }; let job_monitor = async { - monitor - .run_with_signal(shutdown_signal_io()) - .await - .map_err(|e| anyhow::anyhow!("Job monitor error: {e}")) + let result = + monitor.run_with_signal(shutdown_signal_io()).await.context("Job monitor error"); + tracing::info!("Job monitor stopped"); + result }; // Wait for both to complete gracefully (early return on error) - if let Err(e) = tokio::try_join!(web_server, job_monitor) { + if let Err(e) = tokio::try_join!(web_server, jobs_server, job_monitor) { tracing::error!("{e}"); } @@ -163,6 +200,7 @@ fn app(state: AppState, session_store: impl SessionStore + Clone) -> Router { StatusCode::REQUEST_TIMEOUT, Duration::from_secs(120), )) + .layer(NormalizePathLayer::trim_trailing_slash()) .layer(CorsLayer::new().allow_methods([Method::GET]).allow_origin(cors::Any)) .layer( SessionManagerLayer::new(session_store) @@ -172,13 +210,27 @@ fn app(state: AppState, session_store: impl SessionStore + Clone) -> Router { ) .layer(middleware::from_fn(csp_middleware)) .compression(); - let router = build_router(); + let router = build_router().with_state(state); #[cfg(debug_assertions)] let router = router.layer(tower_livereload::LiveReloadLayer::new()); - router.layer(middleware).with_state(state) + router.layer(middleware) } -async fn shutdown_signal() { shutdown_signal_io().await.ok(); } +fn jobs_app(jobs: &JobStorage, broadcaster: Arc>) -> Router { + let middleware = + ServiceBuilder::new().layer(NormalizePathLayer::trim_trailing_slash()).compression(); + let api = ApiBuilder::new(Router::new()) + .register(jobs.workflow_run()) + .register(jobs.refresh_project()) + .build(); + Router::new() + .nest("/api/v1", api) + .fallback_service(ServeUI::new()) + .layer(Extension(broadcaster)) + .layer(middleware) +} + +async fn shutdown_signal() { shutdown_signal_io().await.unwrap() } /// Shutdown signal that returns io::Result for apalis compatibility. async fn shutdown_signal_io() -> std::io::Result<()> {