diff --git a/Cargo.lock b/Cargo.lock index 314944c..35917c0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4717,6 +4717,7 @@ dependencies = [ "thiserror 2.0.19", "time", "tokio", + "tower-http 0.7.0", "tracing", "tracing-subscriber", "uuid", @@ -5332,6 +5333,31 @@ dependencies = [ "tracing", ] +[[package]] +name = "tower-http" +version = "0.7.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b11f75e912b0c2be01b63d8cf8057b8c3f97cf34abb3d431a3a4c8675498e233" +dependencies = [ + "bitflags 2.11.0", + "bytes", + "futures-core", + "futures-util", + "http 1.4.0", + "http-body 1.0.1", + "http-body-util", + "http-range-header 0.4.2", + "httpdate", + "mime", + "mime_guess", + "percent-encoding", + "pin-project-lite", + "tokio", + "tokio-util", + "tower-layer", + "tower-service", +] + [[package]] name = "tower-layer" version = "0.3.3" diff --git a/afterglow/examples/axum_demo.rs b/afterglow/examples/axum_demo.rs index c6b1688..a41fa6d 100644 --- a/afterglow/examples/axum_demo.rs +++ b/afterglow/examples/axum_demo.rs @@ -91,7 +91,10 @@ const FOOTER: Node = html! { async fn index() -> Response { let page = html! { - "afterglow demo" + + "afterglow demo" + +

"Dashboard"

diff --git a/shlepper/Cargo.toml b/shlepper/Cargo.toml index dd62038..5a2f669 100644 --- a/shlepper/Cargo.toml +++ b/shlepper/Cargo.toml @@ -19,6 +19,7 @@ serde = { version = "1.0.229", features = ["derive"] } thiserror = "2.0.19" time = { version = "0.3.53", features = ["serde"] } tokio = { version = "1.53.0", features = ["full"] } +tower-http = { version = "0.7.0", features = ["fs"] } tracing = "0.1.44" tracing-subscriber = { version = "0.3.23", features = ["env-filter"] } uuid = "1.24.0" diff --git a/shlepper/src/cookie.rs b/shlepper/src/cookie.rs index 9ffb168..b58ed53 100644 --- a/shlepper/src/cookie.rs +++ b/shlepper/src/cookie.rs @@ -24,7 +24,7 @@ impl From<[u8; Key::LENGTH]> for Key { impl FromRef for AxumKey { fn from_ref(state: &State) -> Self { - state.cookie_key.clone() + state.cookie_key.0.clone() } } diff --git a/shlepper/src/docker.rs b/shlepper/src/docker.rs new file mode 100644 index 0000000..fbb0ac5 --- /dev/null +++ b/shlepper/src/docker.rs @@ -0,0 +1,16 @@ +use bollard::{API_DEFAULT_VERSION, Docker}; + +pub(super) fn set_up() -> Result { + if cfg!(debug_assertions) { + // I think this is docker desktop specific + Docker::connect_with_socket( + "/Users/claas/.docker/run/docker.sock", + 120, + API_DEFAULT_VERSION, + ) + // Docker::connect_with_unix("unix:///var/run/docker.sock", 120, API_DEFAULT_VERSION) + // Docker::connect_with_unix_defaults() + } else { + Docker::connect_with_defaults() + } +} diff --git a/shlepper/src/error.rs b/shlepper/src/error.rs index f5540e5..a35e75b 100644 --- a/shlepper/src/error.rs +++ b/shlepper/src/error.rs @@ -1,19 +1,13 @@ -use crate::{database, secret}; +use std::io; + +use crate::state; #[derive(thiserror::Error, Debug)] pub(super) enum Error { - #[error("Error setting up secrets")] - SecretError(#[from] secret::Error), - #[error("Error setting up docker")] - DockerError(#[from] bollard::errors::Error), - #[error("Error decoding cookie key")] - CookieDecodeError(base64::DecodeError), - #[error("Bad cookie key length")] - BadCookieKeyLength { expected: usize, actual: usize }, - #[error("Error reading database URL: {0}")] - DatabaseUrlError(#[from] std::env::VarError), - #[error("Bad database key encoding: {0}")] - BadDatabaseKey(base64::DecodeError), - #[error("Error initializing database: {0}")] - DatabaseError(#[from] database::InitializeError), + #[error("Error initializing state: {0}")] + InitializeState(#[from] state::InitializeError), + #[error("Error setting up TCP listener: {0}")] + TcpListener(io::Error), + #[error("Error serving axum app: {0}")] + AxumServe(io::Error), } diff --git a/shlepper/src/garbage_collector.rs b/shlepper/src/garbage_collector.rs new file mode 100644 index 0000000..32853b0 --- /dev/null +++ b/shlepper/src/garbage_collector.rs @@ -0,0 +1,75 @@ +use std::{collections::HashMap, time::Duration}; + +use bollard::query_parameters::ListContainersOptions; + +use crate::state::State; + +mod label { + /// Tag to mark containers managed by tugboat + pub(super) const TAG: &str = "moe.cla.tugboat.tugged"; +} + +/// Runs forever and cleans up expired app data +pub(crate) async fn start(state: State) { + // It is not important that it cleans exactly, but it is important that it happens regularly + // Duration from minutes is experimental currently + let mut interval = tokio::time::interval(Duration::from_secs(24 * 60 * 60)); + let filters = Some(HashMap::from([( + "label".to_owned(), + vec![label::TAG.to_owned()], + )])); + loop { + interval.tick().await; + // Clean up dead containers + let result = state + .docker + .list_containers(Some(ListContainersOptions { + all: true, + filters: filters.clone(), + ..Default::default() + })) + .await; + + let container_ids: Vec = match result { + Ok(containers) => containers + .into_iter() + .filter_map(|container| container.id) + .collect(), + Err(error) => { + tracing::error!("Error listing containers for database cleanup: {:?}", error); + continue; + } + }; + + if container_ids.is_empty() { + continue; + } + + if container_ids.len() >= i16::MAX as usize { + tracing::error!( + "Expected containers to not exceed parameter limit, got {}", + container_ids.len() + ); + continue; + } + + let query = format!( + "DELETE FROM tokens WHERE container_id NOT IN ({})", + (1..=container_ids.len()) + .map(|index| format!("?{index}")) + .collect::>() + .join(", ") + ); + + match state.connection.execute(&query, container_ids).await { + Ok(deleted_rows) => { + tracing::info!( + "Cleaned up {deleted_rows} containers from database that no longer exist" + ) + } + Err(error) => { + tracing::error!("Error cleaning up containers from database: {:?}", error) + } + }; + } +} diff --git a/shlepper/src/main.rs b/shlepper/src/main.rs index 427fe8b..130af36 100644 --- a/shlepper/src/main.rs +++ b/shlepper/src/main.rs @@ -1,13 +1,20 @@ mod cookie; mod database; +mod docker; mod error; +mod garbage_collector; +mod public_directory; mod secret; +mod shutdown_signal; mod state; -use base64::{Engine, engine::general_purpose::URL_SAFE_NO_PAD}; +use std::net::Ipv4Addr; + +use axum::Router; +use tokio::net::TcpListener; use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt}; -use crate::{cookie::Key, error::Error}; +use crate::{error::Error, state::State}; #[tokio::main] async fn main() -> Result<(), Error> { @@ -27,26 +34,30 @@ async fn main() -> Result<(), Error> { #[cfg(debug_assertions)] dotenvy::dotenv().expect("Expected to load .env file in development"); - let secrets = secret::setup().await.inspect_err(|error| { - tracing::error!("Error setting up secrets {}", error); - })?; - - let url = std::env::var("LIBSQL_URL").map_err(Error::DatabaseUrlError)?; - let key = URL_SAFE_NO_PAD - .decode(secrets.database_encryption_key.as_ref()) - .map_err(Error::BadDatabaseKey)? - .into(); - - let connection = database::initialize(url, secrets.lib_sql_auth_token.clone(), key).await?; - - let cookie_key: [u8; Key::LENGTH] = URL_SAFE_NO_PAD - .decode(secrets.cookie_signing_secret.as_ref()) - .map_err(Error::CookieDecodeError)? - .try_into() - .map_err(|secret: Vec| Error::BadCookieKeyLength { - expected: Key::LENGTH, - actual: secret.len(), - })?; + let state = State::initialize().await?; + + // Set up background workers + let _handle = tokio::spawn(garbage_collector::start(state.clone())); + + let app = Router::new() + .fallback_service(public_directory::serve()) + .with_state(state); + + let address = if cfg!(debug_assertions) { + Ipv4Addr::LOCALHOST + } else { + Ipv4Addr::UNSPECIFIED + }; + + let listener = TcpListener::bind((address, 3001)) + .await + .map_err(Error::TcpListener)?; + + tracing::info!("listening on http://{}", listener.local_addr().unwrap()); + axum::serve(listener, app) + .with_graceful_shutdown(shutdown_signal::get()) + .await + .map_err(Error::AxumServe)?; Ok(()) } diff --git a/shlepper/src/public_directory.rs b/shlepper/src/public_directory.rs new file mode 100644 index 0000000..6bdebe0 --- /dev/null +++ b/shlepper/src/public_directory.rs @@ -0,0 +1,32 @@ +use std::path::{Path, PathBuf}; + +use tower_http::services::ServeDir; + +fn get_public_path() -> PathBuf { + if cfg!(debug_assertions) { + Path::new(env!("CARGO_MANIFEST_DIR")).join("public") + } else { + let mut path = std::env::current_exe().unwrap_or_else(|error| { + tracing::warn!( + "Could not get current executable path. Will serve static files from relative \"public\" directory. Causing Error: {}", + error + ); + "public".into() + }); + + // We want the directory containing the executable not the executable itself + _ = path.pop(); + + path.join("public") + } +} + +pub(super) fn serve() -> ServeDir { + let public_path = get_public_path(); + tracing::debug!("Serving files from: {}", public_path.display()); + ServeDir::new(public_path) + .precompressed_br() + .precompressed_deflate() + .precompressed_gzip() + .precompressed_zstd() +} diff --git a/shlepper/src/secret.rs b/shlepper/src/secret.rs index 5c40cea..71be6c8 100644 --- a/shlepper/src/secret.rs +++ b/shlepper/src/secret.rs @@ -54,7 +54,7 @@ pub(super) enum Error { #[error("Error logging in to Bitwarden: {0}")] LoginError(#[from] LoginError), #[error("Error fetching secrets by id from Bitwarden Secrets Manager: {0}")] - GetSecretsError(String), + GetSecretsError(#[source] Box), #[error("Error authenticating with Bitwarden")] BwsAuthenticationFailed, #[error("Error loading secret id from environment variables: {0}")] @@ -121,7 +121,7 @@ pub(super) async fn setup() -> Result { .secrets() .get_by_ids(request) .await - .map_err(|error| Error::GetSecretsError(error.to_string()))?; + .map_err(|error| Error::GetSecretsError(Box::new(error)))?; let mut user_secret = None; let mut cookie_signing_secret = None; let mut lib_sql_auth_token = None; diff --git a/shlepper/src/shutdown_signal.rs b/shlepper/src/shutdown_signal.rs new file mode 100644 index 0000000..a31e4c5 --- /dev/null +++ b/shlepper/src/shutdown_signal.rs @@ -0,0 +1,27 @@ +use tokio::signal; + +pub(super) async fn get() { + let control_c = async { + signal::ctrl_c() + .await + .expect("Failed to install Ctrl+C handler"); + }; + + #[cfg(unix)] + let terminate = async { + signal::unix::signal(signal::unix::SignalKind::terminate()) + .expect("Failed to install signal handler") + .recv() + .await; + }; + + #[cfg(not(unix))] + let terminate = std::future::pending::<()>(); + + tokio::select! { + _ = control_c => {}, + _ = terminate => {}, + } + + tracing::info!("Termination requested. Shutting down"); +} diff --git a/shlepper/src/state.rs b/shlepper/src/state.rs index 7eee615..1f116df 100644 --- a/shlepper/src/state.rs +++ b/shlepper/src/state.rs @@ -1,20 +1,80 @@ use std::{collections::HashMap, sync::Arc}; +use base64::{Engine, engine::general_purpose::URL_SAFE_NO_PAD}; use bollard::Docker; use tokio::sync::Mutex; -use crate::secret::Secrets; +use crate::{ + cookie::{self, Key}, + database, docker, + secret::{self, Secrets}, +}; type UpdateLocks = Arc, Arc>>>>; #[derive(Clone)] pub(crate) struct State { - docker: Docker, - secrets: Secrets, + pub(super) docker: Docker, + pub(super) secrets: Secrets, pub(crate) cookie_key: cookie::Key, // Is there a better primitive to have one task exclusively running the update /// Lock to avoid multiple updates at the same time /// Does not lock the docker instance as other tasks are still permitted - connection: libsql::Connection, - update_locks: UpdateLocks, + pub(super) connection: libsql::Connection, + pub(super) update_locks: UpdateLocks, +} + +#[derive(Debug, thiserror::Error)] +pub(super) enum InitializeError { + #[error("Error setting up secrets")] + Secret(#[from] secret::Error), + #[error("Bad database key encoding: {0}")] + BadDatabaseKey(base64::DecodeError), + #[error("Error reading database URL: {0}")] + DatabaseUrlError(#[from] std::env::VarError), + #[error("Error decoding cookie key")] + CookieDecodeError(base64::DecodeError), + #[error("Bad cookie key length")] + BadCookieKeyLength { expected: usize, actual: usize }, + #[error("Error setting up docker")] + DockerError(#[from] bollard::errors::Error), + #[error("Error initializing database: {0}")] + DatabaseError(#[from] database::InitializeError), +} + +impl State { + pub(super) async fn initialize() -> Result { + let secrets = secret::setup().await.inspect_err(|error| { + tracing::error!("Error setting up secrets {}", error); + })?; + + let key = URL_SAFE_NO_PAD + .decode(secrets.database_encryption_key.as_ref()) + .map_err(InitializeError::BadDatabaseKey)? + .into(); + + let url = std::env::var("LIBSQL_URL").map_err(InitializeError::DatabaseUrlError)?; + let connection = database::initialize(url, secrets.lib_sql_auth_token.clone(), key).await?; + + let cookie_key: [u8; Key::LENGTH] = URL_SAFE_NO_PAD + .decode(secrets.cookie_signing_secret.as_ref()) + .map_err(InitializeError::CookieDecodeError)? + .try_into() + .map_err(|secret: Vec| InitializeError::BadCookieKeyLength { + expected: Key::LENGTH, + actual: secret.len(), + })?; + + let cookie_key = cookie::Key::from(cookie_key); + + let docker = docker::set_up()?; + + Ok(Self { + docker, + secrets, + cookie_key, + connection, + update_locks: Arc::default(), + }) + } }