Integrate apalis-board

This commit is contained in:
Luke Street
2026-01-02 20:13:35 -07:00
parent bfc2b7a7e3
commit 7319c8715e
7 changed files with 199 additions and 46 deletions
Generated
+97 -14
View File
@@ -93,9 +93,9 @@ checksum = "a23eb6b1614318a8071c9b2521f36b424b2c83db5eb3a0fead4a6c0809af6e61"
[[package]] [[package]]
name = "apalis" name = "apalis"
version = "1.0.0-beta.2" version = "1.0.0-rc.1"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c6556b89bb5cb40dab6e4d7ee1f2abfef6d372bd3aef300a3faa992bcb13bda8" checksum = "f93be0eb33b912f5e66004d0b756423c285273259068b1c80a71d7842658189b"
dependencies = [ dependencies = [
"apalis-core", "apalis-core",
"futures-util", "futures-util",
@@ -106,10 +106,59 @@ dependencies = [
] ]
[[package]] [[package]]
name = "apalis-core" name = "apalis-board"
version = "1.0.0-beta.2" version = "1.0.0-rc.1"
source = "registry+https://github.com/rust-lang/crates.io-index" 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 = [ dependencies = [
"futures-channel", "futures-channel",
"futures-core", "futures-core",
@@ -118,7 +167,6 @@ dependencies = [
"futures-util", "futures-util",
"pin-project", "pin-project",
"serde", "serde",
"serde_json",
"thiserror 2.0.17", "thiserror 2.0.17",
"tower-layer", "tower-layer",
"tower-service", "tower-service",
@@ -127,9 +175,9 @@ dependencies = [
[[package]] [[package]]
name = "apalis-cron" 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" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "53230aa8a09cdbe2e254d1b59999efe9f1cad09b66d48475f51b3fb38be6b840" checksum = "d82cc92dc7b6c48c29f308379c6095d82738ebe0831d9bfb27b8a1a8988ee981"
dependencies = [ dependencies = [
"apalis-core", "apalis-core",
"chrono", "chrono",
@@ -141,9 +189,9 @@ dependencies = [
[[package]] [[package]]
name = "apalis-sql" 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" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "604ed96ae8ff20d4c6f8533ffc9f8c91c5f99da6bc564de6a37b94829522f9ec" checksum = "5ade5d8faa60e9975b01d3bb1ebc5028589aa4986365eaa4d080d30ed3b5141f"
dependencies = [ dependencies = [
"apalis-core", "apalis-core",
"chrono", "chrono",
@@ -154,10 +202,11 @@ dependencies = [
[[package]] [[package]]
name = "apalis-sqlite" 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" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8be443bddc6ba6e023b8364304def0a6e3428f258b5addf5aa74bcb5688eb3d1" checksum = "fd43020ce13d6cb8c8d8c09657a6691d8490cd1ce0d8bc0f7fef8bf9b23cfe86"
dependencies = [ dependencies = [
"apalis-codec",
"apalis-core", "apalis-core",
"apalis-sql", "apalis-sql",
"chrono", "chrono",
@@ -174,9 +223,9 @@ dependencies = [
[[package]] [[package]]
name = "apalis-workflow" 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" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "afa1869ce395c0b727315a2be39508eb77ef94ed75492e7b43334b13241b1afa" checksum = "bc024da2d5d3ab59cc9fea099a2e2b20de5ff608f2e287abcb73aa45e4966a89"
dependencies = [ dependencies = [
"apalis-core", "apalis-core",
"futures", "futures",
@@ -1104,6 +1153,7 @@ version = "0.1.0"
dependencies = [ dependencies = [
"anyhow", "anyhow",
"apalis", "apalis",
"apalis-codec",
"apalis-cron", "apalis-cron",
"apalis-sqlite", "apalis-sqlite",
"apalis-workflow", "apalis-workflow",
@@ -1125,6 +1175,7 @@ version = "0.1.0"
dependencies = [ dependencies = [
"anyhow", "anyhow",
"apalis", "apalis",
"apalis-board",
"axum", "axum",
"axum_typed_multipart", "axum_typed_multipart",
"bytes", "bytes",
@@ -2160,6 +2211,25 @@ version = "1.12.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e7c5cedc30da3a610cac6b4ba17597bdf7152cf974e8aab3afb3d54455e371c8" 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]] [[package]]
name = "indexmap" name = "indexmap"
version = "2.12.1" version = "2.12.1"
@@ -4925,6 +4995,16 @@ dependencies = [
"tracing-core", "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]] [[package]]
name = "tracing-subscriber" name = "tracing-subscriber"
version = "0.3.22" version = "0.3.22"
@@ -4935,12 +5015,15 @@ dependencies = [
"nu-ansi-term", "nu-ansi-term",
"once_cell", "once_cell",
"regex-automata", "regex-automata",
"serde",
"serde_json",
"sharded-slab", "sharded-slab",
"smallvec", "smallvec",
"thread_local", "thread_local",
"tracing", "tracing",
"tracing-core", "tracing-core",
"tracing-log", "tracing-log",
"tracing-serde",
] ]
[[package]] [[package]]
+5 -3
View File
@@ -18,9 +18,11 @@ rust-version = "1.88"
[workspace.dependencies] [workspace.dependencies]
anyhow = "1.0" anyhow = "1.0"
apalis = "1.0.0-beta.2" apalis = "1.0.0-rc.1"
apalis-cron = "1.0.0-beta.2" apalis-board = { version = "1.0.0-rc.1", features = ["axum"] }
apalis-sqlite = "1.0.0-beta.4" 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"] } axum = { version = "0.8", features = ["macros"] }
futures-util = "0.3.31" futures-util = "0.3.31"
hex = "0.4" hex = "0.4"
+2 -1
View File
@@ -13,6 +13,7 @@ pub struct Config {
#[derive(Debug, Clone, Deserialize, Serialize)] #[derive(Debug, Clone, Deserialize, Serialize)]
pub struct ServerConfig { pub struct ServerConfig {
pub port: u16, pub port: u16,
pub jobs_port: Option<u16>,
#[serde(default)] #[serde(default)]
pub dev_mode: bool, pub dev_mode: bool,
} }
@@ -64,6 +65,6 @@ pub struct WorkerConfig {
impl Default for WorkerConfig { impl Default for WorkerConfig {
fn default() -> Self { 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 }
} }
} }
+2 -1
View File
@@ -6,9 +6,10 @@ publish = false
[dependencies] [dependencies]
anyhow.workspace = true anyhow.workspace = true
apalis.workspace = true apalis-codec.workspace = true
apalis-cron.workspace = true apalis-cron.workspace = true
apalis-sqlite.workspace = true apalis-sqlite.workspace = true
apalis.workspace = true
apalis-workflow = "0.1.0-beta.2" apalis-workflow = "0.1.0-beta.2"
decomp-dev-core = { path = "../core" } decomp-dev-core = { path = "../core" }
decomp-dev-db = { path = "../db" } decomp-dev-db = { path = "../db" }
+16 -3
View File
@@ -10,6 +10,7 @@ use apalis::{
}, },
prelude::*, prelude::*,
}; };
use apalis_codec::json::JsonCodec;
use apalis_sqlite::{CompactType, SqliteStorage, fetcher::SqliteFetcher}; use apalis_sqlite::{CompactType, SqliteStorage, fetcher::SqliteFetcher};
use decomp_dev_core::config::{Config, DbConfig, WorkerConfig}; use decomp_dev_core::config::{Config, DbConfig, WorkerConfig};
use decomp_dev_db::Database; use decomp_dev_db::Database;
@@ -28,7 +29,7 @@ pub struct JobContext {
} }
/// Type alias for the default codec used by SqliteStorage. /// Type alias for the default codec used by SqliteStorage.
type DefaultCodec = json::JsonCodec<CompactType>; type DefaultCodec = JsonCodec<CompactType>;
/// Type alias for workflow run storage. /// Type alias for workflow run storage.
pub type WorkflowRunStorage = SqliteStorage<ProcessWorkflowRunJob, DefaultCodec, SqliteFetcher>; pub type WorkflowRunStorage = SqliteStorage<ProcessWorkflowRunJob, DefaultCodec, SqliteFetcher>;
@@ -55,8 +56,8 @@ impl JobStorage {
SqlitePool::connect(&db.jobs_url).await.context("Failed to connect to database")?; SqlitePool::connect(&db.jobs_url).await.context("Failed to connect to database")?;
SqliteStorage::setup(&pool).await?; SqliteStorage::setup(&pool).await?;
Ok(Arc::new(Self { Ok(Arc::new(Self {
workflow_run: SqliteStorage::new(&pool), workflow_run: create_storage(&pool),
refresh_project: SqliteStorage::new(&pool), refresh_project: create_storage(&pool),
})) }))
} }
@@ -67,6 +68,18 @@ impl JobStorage {
pub fn refresh_project(&self) -> RefreshProjectStorage { self.refresh_project.clone() } pub fn refresh_project(&self) -> RefreshProjectStorage { self.refresh_project.clone() }
} }
fn create_storage<T>(pool: &SqlitePool) -> SqliteStorage<T, DefaultCodec, SqliteFetcher> {
let config = apalis_sqlite::Config::new(std::any::type_name::<T>()).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. /// Create the job monitor with all workers.
pub fn create_monitor( pub fn create_monitor(
storage: Arc<JobStorage>, storage: Arc<JobStorage>,
+1
View File
@@ -8,6 +8,7 @@ default-run = "decomp-dev-web"
[dependencies] [dependencies]
anyhow.workspace = true anyhow.workspace = true
apalis.workspace = true apalis.workspace = true
apalis-board.workspace = true
axum.workspace = true axum.workspace = true
axum_typed_multipart = "0.16" axum_typed_multipart = "0.16"
decomp-dev-auth = { path = "../auth" } decomp-dev-auth = { path = "../auth" }
+76 -24
View File
@@ -9,12 +9,18 @@ use std::{
io::BufReader, io::BufReader,
net::{IpAddr, Ipv4Addr, SocketAddr}, net::{IpAddr, Ipv4Addr, SocketAddr},
str::FromStr, str::FromStr,
sync::Arc, sync::{Arc, Mutex},
time::Duration, time::Duration,
}; };
use anyhow::Context;
use apalis_board::axum::{
framework::{ApiBuilder, RegisterRoute},
sse::{TracingBroadcaster, TracingSubscriber},
ui::ServeUI,
};
use axum::{ use axum::{
Router, Extension, Router,
extract::{ConnectInfo, FromRef}, extract::{ConnectInfo, FromRef},
http::{Method, Request, StatusCode, header}, http::{Method, Request, StatusCode, header},
middleware, middleware,
@@ -26,15 +32,18 @@ use decomp_dev_jobs::{JobContext, JobStorage, create_monitor};
use tokio::{net::TcpListener, signal}; use tokio::{net::TcpListener, signal};
use tower::ServiceBuilder; use tower::ServiceBuilder;
use tower_http::{ use tower_http::{
ServiceBuilderExt, cors, ServiceBuilderExt,
cors::CorsLayer, cors::{self, CorsLayer},
normalize_path::NormalizePathLayer,
timeout::TimeoutLayer, timeout::TimeoutLayer,
trace::{DefaultOnResponse, MakeSpan, TraceLayer}, trace::{DefaultOnResponse, MakeSpan, TraceLayer},
}; };
use tower_sessions::{Expiry, SessionManagerLayer, SessionStore, cookie::SameSite}; use tower_sessions::{Expiry, SessionManagerLayer, SessionStore, cookie::SameSite};
use tower_sessions_sqlx_store::SqliteStore; use tower_sessions_sqlx_store::SqliteStore;
use tracing::{Level, Span}; 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}; use crate::handlers::{build_router, csp::csp_middleware};
@@ -48,13 +57,14 @@ pub struct AppState {
#[tokio::main] #[tokio::main]
async fn main() { async fn main() {
tracing_subscriber::fmt() let broadcaster = TracingBroadcaster::create();
.with_env_filter( let env_filter = EnvFilter::builder()
EnvFilter::builder() // Default to info level
// Default to info level .with_default_directive(LevelFilter::INFO.into())
.with_default_directive(LevelFilter::INFO.into()) .from_env_lossy();
.from_env_lossy(), tracing_subscriber::registry()
) .with(tracing_subscriber::fmt::layer().with_filter(env_filter.clone()))
.with(TracingSubscriber::new(&broadcaster).layer())
.init(); .init();
let config: Arc<Config> = { let config: Arc<Config> = {
@@ -86,6 +96,8 @@ async fn main() {
// Build the router // Build the router
let port = state.config.server.port; 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::<SocketAddr>(); let router = app(state, session_store).into_make_service_with_connect_info::<SocketAddr>();
// Create the listener // Create the listener
@@ -97,7 +109,7 @@ async fn main() {
let fds = libsystemd::activation::receive_descriptors_with_names(false) let fds = libsystemd::activation::receive_descriptors_with_names(false)
.expect("Failed to receive fds"); .expect("Failed to receive fds");
if let Some((fd, name)) = fds.into_iter().next() { 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()) }; 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"); std_listener.set_nonblocking(true).expect("Failed to set non-blocking");
listener = listener =
@@ -108,10 +120,21 @@ async fn main() {
Some(listener) => listener, Some(listener) => listener,
None => { None => {
let addr = SocketAddr::from((Ipv4Addr::UNSPECIFIED, port)); 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") 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")] #[cfg(target_os = "linux")]
{ {
@@ -121,20 +144,34 @@ async fn main() {
// Run both the web server and job monitor concurrently, with graceful shutdown // Run both the web server and job monitor concurrently, with graceful shutdown
let web_server = async { let web_server = async {
axum::serve(listener, router) let result = axum::serve(listener, router)
.with_graceful_shutdown(shutdown_signal()) .with_graceful_shutdown(shutdown_signal())
.await .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 { let job_monitor = async {
monitor let result =
.run_with_signal(shutdown_signal_io()) monitor.run_with_signal(shutdown_signal_io()).await.context("Job monitor error");
.await tracing::info!("Job monitor stopped");
.map_err(|e| anyhow::anyhow!("Job monitor error: {e}")) result
}; };
// Wait for both to complete gracefully (early return on error) // 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}"); tracing::error!("{e}");
} }
@@ -163,6 +200,7 @@ fn app(state: AppState, session_store: impl SessionStore + Clone) -> Router {
StatusCode::REQUEST_TIMEOUT, StatusCode::REQUEST_TIMEOUT,
Duration::from_secs(120), Duration::from_secs(120),
)) ))
.layer(NormalizePathLayer::trim_trailing_slash())
.layer(CorsLayer::new().allow_methods([Method::GET]).allow_origin(cors::Any)) .layer(CorsLayer::new().allow_methods([Method::GET]).allow_origin(cors::Any))
.layer( .layer(
SessionManagerLayer::new(session_store) SessionManagerLayer::new(session_store)
@@ -172,13 +210,27 @@ fn app(state: AppState, session_store: impl SessionStore + Clone) -> Router {
) )
.layer(middleware::from_fn(csp_middleware)) .layer(middleware::from_fn(csp_middleware))
.compression(); .compression();
let router = build_router(); let router = build_router().with_state(state);
#[cfg(debug_assertions)] #[cfg(debug_assertions)]
let router = router.layer(tower_livereload::LiveReloadLayer::new()); 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<Mutex<TracingBroadcaster>>) -> 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. /// Shutdown signal that returns io::Result for apalis compatibility.
async fn shutdown_signal_io() -> std::io::Result<()> { async fn shutdown_signal_io() -> std::io::Result<()> {