From f5ce3570e4d041cfcb208849acf176b9c8b240df Mon Sep 17 00:00:00 2001 From: mbecker20 Date: Sun, 2 Jun 2024 04:14:51 -0700 Subject: [PATCH] execute api returns update immediately --- Cargo.lock | 12 +- Cargo.toml | 2 +- bin/core/src/api/execute/build.rs | 48 +++--- bin/core/src/api/execute/deployment.rs | 109 +++++-------- bin/core/src/api/execute/mod.rs | 75 ++++++--- bin/core/src/api/execute/procedure.rs | 24 +-- bin/core/src/api/execute/repo.rs | 44 +----- bin/core/src/api/execute/server.rs | 33 +--- bin/core/src/api/execute/server_template.rs | 12 +- bin/core/src/helpers/procedure.rs | 163 +++++++++++++++----- bin/core/src/helpers/update.rs | 80 +++++++++- bin/core/src/listener/github.rs | 70 +++++---- bin/migrator/src/legacy/v0/mod.rs | 8 +- client/core/rs/src/entities/mod.rs | 8 +- 14 files changed, 398 insertions(+), 290 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 51be9195f..30842d93a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3088,9 +3088,9 @@ checksum = "388a1df253eca08550bef6c72392cfe7c30914bf41df5269b68cbd6ff8f570a3" [[package]] name = "serde" -version = "1.0.202" +version = "1.0.203" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "226b61a0d411b2ba5ff6d7f73a476ac4f8bb900373459cd00fab8512828ba395" +checksum = "7253ab4de971e72fb7be983802300c30b5a7f0c2e56fab8abfc6a214307c0094" dependencies = [ "serde_derive", ] @@ -3106,9 +3106,9 @@ dependencies = [ [[package]] name = "serde_derive" -version = "1.0.202" +version = "1.0.203" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6048858004bcff69094cd972ed40a32500f153bd3be9f716b2eed2e8217c4838" +checksum = "500cbc0ebeb6f46627f50f3f5811ccf6bf00643be300b4c3eabc0ef55dc5b5ba" dependencies = [ "proc-macro2", "quote", @@ -3210,9 +3210,9 @@ dependencies = [ [[package]] name = "serror" -version = "0.3.4" +version = "0.3.5" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "c6c430d274ef4964c27e7b338fdfed6d400b18ceac80cba827ddd84055f9eca8" +checksum = "d4c2d61f8b059a9ac51ab4f65ea9874cb1b8a6e57f2ebc6d284988f6b877253b" dependencies = [ "anyhow", "axum 0.7.5", diff --git a/Cargo.toml b/Cargo.toml index 0c3b8b27e..8e093bc66 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -21,7 +21,7 @@ logger = { path = "lib/logger" } # MOGH run_command = { version = "0.0.6", features = ["async_tokio"] } -serror = { version = "0.3.4", default-features = false } +serror = { version = "0.3.5", default-features = false } slack = { version = "0.1.0", package = "slack_client_rs" } derive_default_builder = "0.1.8" derive_empty_traits = "0.1.0" diff --git a/bin/core/src/api/execute/build.rs b/bin/core/src/api/execute/build.rs index 7d0bbd043..89e76155e 100644 --- a/bin/core/src/api/execute/build.rs +++ b/bin/core/src/api/execute/build.rs @@ -52,12 +52,14 @@ use crate::{ state::{action_states, db_client, State}, }; -impl Resolve for State { +use crate::helpers::update::init_execution_update; + +impl Resolve for State { #[instrument(name = "RunBuild", skip(self, user))] async fn resolve( &self, RunBuild { build }: RunBuild, - user: User, + (user, mut update): (User, Update), ) -> anyhow::Result { let mut build = resource::get_check_permissions::( &build, @@ -77,8 +79,6 @@ impl Resolve for State { build.config.version.increment(); - let mut update = make_update(&build, Operation::RunBuild, &user); - update.in_progress(); update.version = build.config.version.clone(); let cancel = CancellationToken::new(); @@ -337,12 +337,12 @@ async fn handle_early_return( Ok(update) } -impl Resolve for State { +impl Resolve for State { #[instrument(name = "CancelBuild", skip(self, user))] async fn resolve( &self, CancelBuild { build }: CancelBuild, - user: User, + (user, mut update): (User, Update), ) -> anyhow::Result { let build = resource::get_check_permissions::( &build, @@ -370,9 +370,6 @@ impl Resolve for State { return Err(anyhow!("Build cancel is already in progress")); } - let mut update = - make_update(&build, Operation::CancelBuild, &user); - update.push_simple_log( "cancel triggered", "the build cancel has been triggered", @@ -564,16 +561,26 @@ async fn handle_post_build_redeploy(build_id: &str) { let state = get_deployment_state(&deployment).await.unwrap_or_default(); if state == DeploymentState::Running { - let res = State - .resolve( - Deploy { - deployment: deployment.id.clone(), - stop_signal: None, - stop_time: None, - }, - auto_redeploy_user().to_owned(), - ) - .await; + let req = super::ExecuteRequest::Deploy(Deploy { + deployment: deployment.id.clone(), + stop_signal: None, + stop_time: None, + }); + let user = auto_redeploy_user().to_owned(); + let res = async { + let update = init_execution_update(&req, &user).await?; + State + .resolve( + Deploy { + deployment: deployment.id.clone(), + stop_signal: None, + stop_time: None, + }, + (user, update), + ) + .await + } + .await; Some((deployment.id.clone(), res)) } else { None @@ -592,7 +599,8 @@ async fn handle_post_build_redeploy(build_id: &str) { let (id, res) = res.unwrap(); match res { Ok(_) => redeploys.push(id), - Err(e) => redeploy_failures.push(format!("{id}: {e:#?}")), + Err(e) => redeploy_failures + .push(format!("{id}: {}", serialize_error_pretty(&e))), } } } diff --git a/bin/core/src/api/execute/deployment.rs b/bin/core/src/api/execute/deployment.rs index cbeabd4f9..a9ad2d144 100644 --- a/bin/core/src/api/execute/deployment.rs +++ b/bin/core/src/api/execute/deployment.rs @@ -7,12 +7,12 @@ use monitor_client::{ entities::{ build::Build, deployment::{Deployment, DeploymentImage}, - get_image_name, monitor_timestamp, + get_image_name, permission::PermissionLevel, server::ServerState, - update::{Log, ResourceTarget, Update, UpdateStatus}, + update::{Log, Update}, user::User, - Operation, Version, + Version, }, }; use mungos::{find::find_collect, mongodb::bson::doc}; @@ -25,14 +25,16 @@ use crate::{ helpers::{ periphery_client, query::{get_global_variables, get_server_with_status}, - update::{add_update, make_update, update_update}, + update::update_update, }, monitor::update_cache_for_server, resource, state::{action_states, db_client, State}, }; -impl Resolve for State { +use crate::helpers::update::init_execution_update; + +impl Resolve for State { #[instrument(name = "Deploy", skip(self, user))] async fn resolve( &self, @@ -41,7 +43,7 @@ impl Resolve for State { stop_signal, stop_time, }: Deploy, - user: User, + (user, mut update): (User, Update), ) -> anyhow::Result { let mut deployment = resource::get_check_permissions::( @@ -129,11 +131,6 @@ impl Resolve for State { env.value = res; } - let mut update = - make_update(&deployment, Operation::DeployContainer, &user); - update.in_progress(); - update.version = version; - // Show which variables were interpolated if !global_replacers.is_empty() { update.push_simple_log( @@ -156,7 +153,8 @@ impl Resolve for State { ); } - update.id = add_update(update.clone()).await?; + update.version = version; + update_update(update.clone()).await?; let docker_token = core_config .docker_accounts @@ -191,12 +189,12 @@ impl Resolve for State { } } -impl Resolve for State { +impl Resolve for State { #[instrument(name = "StartContainer", skip(self, user))] async fn resolve( &self, StartContainer { deployment }: StartContainer, - user: User, + (user, mut update): (User, Update), ) -> anyhow::Result { let deployment = resource::get_check_permissions::( &deployment, @@ -230,20 +228,6 @@ impl Resolve for State { let periphery = periphery_client(&server)?; - let start_ts = monitor_timestamp(); - - let mut update = Update { - target: ResourceTarget::Deployment(deployment.id.clone()), - operation: Operation::StartContainer, - start_ts, - status: UpdateStatus::InProgress, - success: true, - operator: user.id.clone(), - ..Default::default() - }; - - update.id = add_update(update.clone()).await?; - let log = match periphery .request(api::container::StartContainer { name: deployment.name.clone(), @@ -257,15 +241,15 @@ impl Resolve for State { }; update.logs.push(log); - update.finalize(); update_cache_for_server(&server).await; + update.finalize(); update_update(update.clone()).await?; Ok(update) } } -impl Resolve for State { +impl Resolve for State { #[instrument(name = "StopContainer", skip(self, user))] async fn resolve( &self, @@ -274,7 +258,7 @@ impl Resolve for State { signal, time, }: StopContainer, - user: User, + (user, mut update): (User, Update), ) -> anyhow::Result { let deployment = resource::get_check_permissions::( &deployment, @@ -308,11 +292,6 @@ impl Resolve for State { let periphery = periphery_client(&server)?; - let mut update = - make_update(&deployment, Operation::StopContainer, &user); - - update.id = add_update(update.clone()).await?; - let log = match periphery .request(api::container::StopContainer { name: deployment.name.clone(), @@ -332,20 +311,20 @@ impl Resolve for State { }; update.logs.push(log); - update.finalize(); update_cache_for_server(&server).await; + update.finalize(); update_update(update.clone()).await?; Ok(update) } } -impl Resolve for State { +impl Resolve for State { #[instrument(name = "StopAllContainers", skip(self, user))] async fn resolve( &self, StopAllContainers { server }: StopAllContainers, - user: User, + (user, mut update): (User, Update), ) -> anyhow::Result { let (server, status) = get_server_with_status(&server).await?; if status != ServerState::Ok { @@ -375,23 +354,27 @@ impl Resolve for State { .await .context("failed to find deployments on server")?; - let mut update = - make_update(&server, Operation::StopAllContainers, &user); - update.in_progress(); - update.id = add_update(update.clone()).await?; - let futures = deployments.iter().map(|deployment| async { + let req = super::ExecuteRequest::StopContainer(StopContainer { + deployment: deployment.id.clone(), + signal: None, + time: None, + }); ( - self - .resolve( - StopContainer { - deployment: deployment.id.clone(), - signal: None, - time: None, - }, - user.clone(), - ) - .await, + async { + let update = init_execution_update(&req, &user).await?; + State + .resolve( + StopContainer { + deployment: deployment.id.clone(), + signal: None, + time: None, + }, + (user.clone(), update), + ) + .await + } + .await, deployment.name.clone(), deployment.id.clone(), ) @@ -422,7 +405,7 @@ impl Resolve for State { } } -impl Resolve for State { +impl Resolve for State { #[instrument(name = "RemoveContainer", skip(self, user))] async fn resolve( &self, @@ -431,7 +414,7 @@ impl Resolve for State { signal, time, }: RemoveContainer, - user: User, + (user, mut update): (User, Update), ) -> anyhow::Result { let deployment = resource::get_check_permissions::( &deployment, @@ -465,20 +448,6 @@ impl Resolve for State { let periphery = periphery_client(&server)?; - let start_ts = monitor_timestamp(); - - let mut update = Update { - target: ResourceTarget::Deployment(deployment.id.clone()), - operation: Operation::RemoveContainer, - start_ts, - status: UpdateStatus::InProgress, - success: true, - operator: user.id.clone(), - ..Default::default() - }; - - update.id = add_update(update.clone()).await?; - let log = match periphery .request(api::container::RemoveContainer { name: deployment.name.clone(), diff --git a/bin/core/src/api/execute/mod.rs b/bin/core/src/api/execute/mod.rs index b21f91d39..c82b9e1ca 100644 --- a/bin/core/src/api/execute/mod.rs +++ b/bin/core/src/api/execute/mod.rs @@ -1,16 +1,22 @@ use std::time::Instant; -use anyhow::{anyhow, Context}; +use anyhow::anyhow; use axum::{middleware, routing::post, Extension, Router}; -use axum_extra::{headers::ContentType, TypedHeader}; -use monitor_client::{api::execute::*, entities::user::User}; +use monitor_client::{ + api::execute::*, + entities::{update::Update, user::User}, +}; use resolver_api::{derive::Resolver, Resolver}; use serde::{Deserialize, Serialize}; -use serror::Json; +use serror::{serialize_error_pretty, Json}; use typeshare::typeshare; use uuid::Uuid; -use crate::{auth::auth_request, state::State}; +use crate::{ + auth::auth_request, + helpers::update::{init_execution_update, update_update}, + state::State, +}; mod build; mod deployment; @@ -22,9 +28,9 @@ mod server_template; #[typeshare] #[derive(Serialize, Deserialize, Debug, Clone, Resolver)] #[resolver_target(State)] -#[resolver_args(User)] +#[resolver_args((User, Update))] #[serde(tag = "type", content = "params")] -enum ExecuteRequest { +pub enum ExecuteRequest { // ==== SERVER ==== PruneContainers(PruneContainers), PruneImages(PruneImages), @@ -61,18 +67,37 @@ pub fn router() -> Router { async fn handler( Extension(user): Extension, Json(request): Json, -) -> serror::Result<(TypedHeader, String)> { +) -> serror::Result> { let req_id = Uuid::new_v4(); - let res = tokio::spawn(task(req_id, request, user)) - .await - .context("failure in spawned execute task"); + let update = init_execution_update(&request, &user).await?; - if let Err(e) = &res { - warn!("/execute request {req_id} spawn error: {e:#}",); - } + let handle = + tokio::spawn(task(req_id, request, user, update.clone())); - Ok((TypedHeader(ContentType::json()), res??)) + tokio::spawn({ + let mut update = update.clone(); + async move { + match handle.await { + Ok(Err(e)) => { + warn!("/execute request {req_id} task error: {e:#}",); + update + .push_error_log("task error", serialize_error_pretty(&e)); + update.finalize(); + let _ = update_update(update).await; + } + Err(e) => { + warn!("/execute request {req_id} spawn error: {e:?}",); + update.push_error_log("spawn error", format!("{e:#?}")); + update.finalize(); + let _ = update_update(update).await; + } + _ => {} + } + } + }); + + Ok(Json(update)) } #[instrument(name = "ExecuteRequest", skip(user))] @@ -80,6 +105,7 @@ async fn task( req_id: Uuid, request: ExecuteRequest, user: User, + update: Update, ) -> anyhow::Result { info!( "/execute request {req_id} | user: {} ({})", @@ -87,16 +113,15 @@ async fn task( ); let timer = Instant::now(); - let res = - State - .resolve_request(request, user) - .await - .map_err(|e| match e { - resolver_api::Error::Serialization(e) => { - anyhow!("{e:?}").context("response serialization error") - } - resolver_api::Error::Inner(e) => e, - }); + let res = State + .resolve_request(request, (user, update)) + .await + .map_err(|e| match e { + resolver_api::Error::Serialization(e) => { + anyhow!("{e:?}").context("response serialization error") + } + resolver_api::Error::Inner(e) => e, + }); if let Err(e) = &res { warn!("/execute request {req_id} error: {e:#}"); diff --git a/bin/core/src/api/execute/procedure.rs b/bin/core/src/api/execute/procedure.rs index f089208c8..3f5488d52 100644 --- a/bin/core/src/api/execute/procedure.rs +++ b/bin/core/src/api/execute/procedure.rs @@ -4,7 +4,7 @@ use monitor_client::{ api::execute::RunProcedure, entities::{ permission::PermissionLevel, procedure::Procedure, - update::Update, user::User, Operation, + update::Update, user::User, }, }; use mungos::{by_id::update_one_by_id, mongodb::bson::to_document}; @@ -13,28 +13,26 @@ use serror::serialize_error_pretty; use tokio::sync::Mutex; use crate::{ - helpers::{ - procedure::execute_procedure, - update::{add_update, make_update, update_update}, - }, + helpers::{procedure::execute_procedure, update::update_update}, resource::{self, refresh_procedure_state_cache}, state::{action_states, db_client, State}, }; -impl Resolve for State { +impl Resolve for State { #[instrument(name = "RunProcedure", skip(self, user))] async fn resolve( &self, RunProcedure { procedure }: RunProcedure, - user: User, + (user, update): (User, Update), ) -> anyhow::Result { - resolve_inner(procedure, user).await + resolve_inner(procedure, user, update).await } } fn resolve_inner( procedure: String, user: User, + update: Update, ) -> Pin< Box< dyn std::future::Future> + Send, @@ -59,16 +57,6 @@ fn resolve_inner( let _action_guard = action_state.update(|state| state.running = true)?; - let mut update = - make_update(&procedure, Operation::RunProcedure, &user); - update.in_progress(); - update.push_simple_log( - "execute procedure", - format!("Executing procedure: {}", procedure.name), - ); - - update.id = add_update(update.clone()).await?; - let update = Mutex::new(update); let res = execute_procedure(&procedure, &update).await; diff --git a/bin/core/src/api/execute/repo.rs b/bin/core/src/api/execute/repo.rs index 136185f95..44172c154 100644 --- a/bin/core/src/api/execute/repo.rs +++ b/bin/core/src/api/execute/repo.rs @@ -6,9 +6,8 @@ use monitor_client::{ permission::PermissionLevel, repo::Repo, server::Server, - update::{Log, ResourceTarget, Update, UpdateStatus}, + update::{Log, Update}, user::User, - Operation, }, }; use mungos::{ @@ -21,20 +20,17 @@ use serror::serialize_error_pretty; use crate::{ config::core_config, - helpers::{ - periphery_client, - update::{add_update, update_update}, - }, + helpers::{periphery_client, update::update_update}, resource::{self, refresh_repo_state_cache}, state::{action_states, db_client, State}, }; -impl Resolve for State { +impl Resolve for State { #[instrument(name = "CloneRepo", skip(self, user))] async fn resolve( &self, CloneRepo { repo }: CloneRepo, - user: User, + (user, mut update): (User, Update), ) -> anyhow::Result { let repo = resource::get_check_permissions::( &repo, @@ -61,20 +57,6 @@ impl Resolve for State { let periphery = periphery_client(&server)?; - let start_ts = monitor_timestamp(); - - let mut update = Update { - operation: Operation::CloneRepo, - target: ResourceTarget::Repo(repo.id.clone()), - start_ts, - status: UpdateStatus::InProgress, - operator: user.id.clone(), - success: true, - ..Default::default() - }; - - update.id = add_update(update.clone()).await?; - let github_token = core_config() .github_accounts .get(&repo.config.github_account) @@ -104,12 +86,12 @@ impl Resolve for State { } } -impl Resolve for State { +impl Resolve for State { #[instrument(name = "PullRepo", skip(self, user))] async fn resolve( &self, PullRepo { repo }: PullRepo, - user: User, + (user, mut update): (User, Update), ) -> anyhow::Result { let repo = resource::get_check_permissions::( &repo, @@ -136,20 +118,6 @@ impl Resolve for State { let periphery = periphery_client(&server)?; - let start_ts = monitor_timestamp(); - - let mut update = Update { - operation: Operation::PullRepo, - target: ResourceTarget::Repo(repo.id.clone()), - start_ts, - status: UpdateStatus::InProgress, - operator: user.id.clone(), - success: true, - ..Default::default() - }; - - update.id = add_update(update.clone()).await?; - let logs = match periphery .request(api::git::PullRepo { name: repo.name.clone(), diff --git a/bin/core/src/api/execute/server.rs b/bin/core/src/api/execute/server.rs index 42b639585..a7468259a 100644 --- a/bin/core/src/api/execute/server.rs +++ b/bin/core/src/api/execute/server.rs @@ -7,7 +7,6 @@ use monitor_client::{ server::Server, update::{Log, Update, UpdateStatus}, user::User, - Operation, }, }; use periphery_client::api; @@ -15,20 +14,17 @@ use resolver_api::Resolve; use serror::serialize_error_pretty; use crate::{ - helpers::{ - periphery_client, - update::{add_update, make_update, update_update}, - }, + helpers::{periphery_client, update::update_update}, resource, state::{action_states, State}, }; -impl Resolve for State { +impl Resolve for State { #[instrument(name = "PruneContainers", skip(self, user))] async fn resolve( &self, PruneContainers { server }: PruneContainers, - user: User, + (user, mut update): (User, Update), ) -> anyhow::Result { let server = resource::get_check_permissions::( &server, @@ -50,11 +46,6 @@ impl Resolve for State { let periphery = periphery_client(&server)?; - let mut update = - make_update(&server, Operation::PruneContainersServer, &user); - update.in_progress(); - update.id = add_update(update.clone()).await?; - let log = match periphery .request(api::container::PruneContainers {}) .await @@ -79,12 +70,12 @@ impl Resolve for State { } } -impl Resolve for State { +impl Resolve for State { #[instrument(name = "PruneNetworks", skip(self, user))] async fn resolve( &self, PruneNetworks { server }: PruneNetworks, - user: User, + (user, mut update): (User, Update), ) -> anyhow::Result { let server = resource::get_check_permissions::( &server, @@ -106,11 +97,6 @@ impl Resolve for State { let periphery = periphery_client(&server)?; - let mut update = - make_update(&server, Operation::PruneNetworksServer, &user); - update.in_progress(); - update.id = add_update(update.clone()).await?; - let log = match periphery .request(api::network::PruneNetworks {}) .await @@ -135,12 +121,12 @@ impl Resolve for State { } } -impl Resolve for State { +impl Resolve for State { #[instrument(name = "PruneImages", skip(self, user))] async fn resolve( &self, PruneImages { server }: PruneImages, - user: User, + (user, mut update): (User, Update), ) -> anyhow::Result { let server = resource::get_check_permissions::( &server, @@ -162,11 +148,6 @@ impl Resolve for State { let periphery = periphery_client(&server)?; - let mut update = - make_update(&server, Operation::PruneImagesServer, &user); - update.in_progress(); - update.id = add_update(update.clone()).await?; - let log = match periphery.request(api::build::PruneImages {}).await { Ok(log) => log, diff --git a/bin/core/src/api/execute/server_template.rs b/bin/core/src/api/execute/server_template.rs index c18f154c9..ef1c19f3c 100644 --- a/bin/core/src/api/execute/server_template.rs +++ b/bin/core/src/api/execute/server_template.rs @@ -7,7 +7,6 @@ use monitor_client::{ server_template::{ServerTemplate, ServerTemplateConfig}, update::Update, user::User, - Operation, }, }; use mungos::mongodb::bson::doc; @@ -16,12 +15,12 @@ use serror::serialize_error_pretty; use crate::{ cloud::{aws::launch_ec2_instance, hetzner::launch_hetzner_server}, - helpers::update::{add_update, make_update, update_update}, + helpers::update::update_update, resource, state::{db_client, State}, }; -impl Resolve for State { +impl Resolve for State { #[instrument(name = "LaunchServer", skip(self, user))] async fn resolve( &self, @@ -29,7 +28,7 @@ impl Resolve for State { name, server_template, }: LaunchServer, - user: User, + (user, mut update): (User, Update), ) -> anyhow::Result { // validate name isn't already taken by another server if db_client() @@ -55,14 +54,11 @@ impl Resolve for State { ) .await?; - let mut update = - make_update(&template, Operation::LaunchServer, &user); - update.in_progress(); update.push_simple_log( "launching server", format!("{:#?}", template.config), ); - update.id = add_update(update.clone()).await?; + update_update(update.clone()).await?; let config = match template.config { ServerTemplateConfig::Aws(config) => { diff --git a/bin/core/src/helpers/procedure.rs b/bin/core/src/helpers/procedure.rs index c0dd4f803..facda1e4b 100644 --- a/bin/core/src/helpers/procedure.rs +++ b/bin/core/src/helpers/procedure.rs @@ -12,9 +12,9 @@ use monitor_client::{ use resolver_api::Resolve; use tokio::sync::Mutex; -use crate::state::State; +use crate::{api::execute::ExecuteRequest, state::State}; -use super::update::update_update; +use super::update::{init_execution_update, update_update}; #[instrument] pub async fn execute_procedure( @@ -79,58 +79,137 @@ async fn execute_execution( if req.procedure == parent_id || req.procedure == parent_name { return Err(anyhow!("Self referential procedure detected")); } + let req = ExecuteRequest::RunProcedure(req); + let update = init_execution_update(&req, &user).await?; + let ExecuteRequest::RunProcedure(req) = req else { + unreachable!() + }; State - .resolve(req, user) + .resolve(req, (user, update)) .await .context("failed at RunProcedure")? } - Execution::RunBuild(req) => State - .resolve(req, user) - .await - .context("failed at RunBuild")?, - Execution::Deploy(req) => { - State.resolve(req, user).await.context("failed at Deploy")? - } - Execution::StartContainer(req) => State - .resolve(req, user) - .await - .context("failed at StartContainer")?, - Execution::StopContainer(req) => { + Execution::RunBuild(req) => { + let req = ExecuteRequest::RunBuild(req); + let update = init_execution_update(&req, &user).await?; + let ExecuteRequest::RunBuild(req) = req else { + unreachable!() + }; State - .resolve(req, user) + .resolve(req, (user, update)) + .await + .context("failed at RunBuild")? + } + Execution::Deploy(req) => { + let req = ExecuteRequest::Deploy(req); + let update = init_execution_update(&req, &user).await?; + let ExecuteRequest::Deploy(req) = req else { + unreachable!() + }; + State + .resolve(req, (user, update)) + .await + .context("failed at Deploy")? + } + Execution::StartContainer(req) => { + let req = ExecuteRequest::StartContainer(req); + let update = init_execution_update(&req, &user).await?; + let ExecuteRequest::StartContainer(req) = req else { + unreachable!() + }; + State + .resolve(req, (user, update)) + .await + .context("failed at StartContainer")? + } + Execution::StopContainer(req) => { + let req = ExecuteRequest::StopContainer(req); + let update = init_execution_update(&req, &user).await?; + let ExecuteRequest::StopContainer(req) = req else { + unreachable!() + }; + State + .resolve(req, (user, update)) .await .context("failed at StopContainer")? } - Execution::StopAllContainers(req) => State - .resolve(req, user) - .await - .context("failed at StopAllContainers")?, - Execution::RemoveContainer(req) => State - .resolve(req, user) - .await - .context("failed at RemoveContainer")?, - Execution::CloneRepo(req) => State - .resolve(req, user) - .await - .context("failed at CloneRepo")?, - Execution::PullRepo(req) => State - .resolve(req, user) - .await - .context("failed at PullRepo")?, - Execution::PruneNetworks(req) => { + Execution::StopAllContainers(req) => { + let req = ExecuteRequest::StopAllContainers(req); + let update = init_execution_update(&req, &user).await?; + let ExecuteRequest::StopAllContainers(req) = req else { + unreachable!() + }; State - .resolve(req, user) + .resolve(req, (user, update)) + .await + .context("failed at StopAllContainers")? + } + Execution::RemoveContainer(req) => { + let req = ExecuteRequest::RemoveContainer(req); + let update = init_execution_update(&req, &user).await?; + let ExecuteRequest::RemoveContainer(req) = req else { + unreachable!() + }; + State + .resolve(req, (user, update)) + .await + .context("failed at RemoveContainer")? + } + Execution::CloneRepo(req) => { + let req = ExecuteRequest::CloneRepo(req); + let update = init_execution_update(&req, &user).await?; + let ExecuteRequest::CloneRepo(req) = req else { + unreachable!() + }; + State + .resolve(req, (user, update)) + .await + .context("failed at CloneRepo")? + } + Execution::PullRepo(req) => { + let req = ExecuteRequest::PullRepo(req); + let update = init_execution_update(&req, &user).await?; + let ExecuteRequest::PullRepo(req) = req else { + unreachable!() + }; + State + .resolve(req, (user, update)) + .await + .context("failed at PullRepo")? + } + Execution::PruneNetworks(req) => { + let req = ExecuteRequest::PruneNetworks(req); + let update = init_execution_update(&req, &user).await?; + let ExecuteRequest::PruneNetworks(req) = req else { + unreachable!() + }; + State + .resolve(req, (user, update)) .await .context("failed at PruneNetworks")? } - Execution::PruneImages(req) => State - .resolve(req, user) - .await - .context("failed at PruneImages")?, - Execution::PruneContainers(req) => State - .resolve(req, user) - .await - .context("failed at PruneContainers")?, + Execution::PruneImages(req) => { + let req = ExecuteRequest::PruneImages(req); + let update = init_execution_update(&req, &user).await?; + let ExecuteRequest::PruneImages(req) = req else { + unreachable!() + }; + State + .resolve(req, (user, update)) + .await + .context("failed at PruneImages")? + } + Execution::PruneContainers(req) => { + let req = ExecuteRequest::PruneContainers(req); + let update = init_execution_update(&req, &user).await?; + let ExecuteRequest::PruneContainers(req) = req else { + unreachable!() + }; + State + .resolve(req, (user, update)) + .await + .context("failed at PruneContainers")? + } }; if update.success { Ok(()) diff --git a/bin/core/src/helpers/update.rs b/bin/core/src/helpers/update.rs index 131b45c07..432fad149 100644 --- a/bin/core/src/helpers/update.rs +++ b/bin/core/src/helpers/update.rs @@ -10,7 +10,7 @@ use mungos::{ mongodb::bson::to_document, }; -use crate::state::db_client; +use crate::{api::execute::ExecuteRequest, state::db_client}; use super::channel::update_channel; @@ -94,3 +94,81 @@ async fn send_update(update: UpdateListItem) -> anyhow::Result<()> { update_channel().sender.lock().await.send(update)?; Ok(()) } + +pub async fn init_execution_update( + request: &ExecuteRequest, + user: &User, +) -> anyhow::Result { + let (operation, target) = match &request { + // Server + ExecuteRequest::PruneContainers(data) => ( + Operation::PruneImages, + ResourceTarget::Server(data.server.clone()), + ), + ExecuteRequest::PruneImages(data) => ( + Operation::PruneImages, + ResourceTarget::Server(data.server.clone()), + ), + ExecuteRequest::PruneNetworks(data) => ( + Operation::PruneNetworks, + ResourceTarget::Server(data.server.clone()), + ), + ExecuteRequest::StopAllContainers(data) => ( + Operation::StopAllContainers, + ResourceTarget::Server(data.server.clone()), + ), + + // Deployment + ExecuteRequest::Deploy(data) => ( + Operation::Deploy, + ResourceTarget::Deployment(data.deployment.clone()), + ), + ExecuteRequest::StartContainer(data) => ( + Operation::StartContainer, + ResourceTarget::Deployment(data.deployment.clone()), + ), + ExecuteRequest::StopContainer(data) => ( + Operation::StopContainer, + ResourceTarget::Deployment(data.deployment.clone()), + ), + ExecuteRequest::RemoveContainer(data) => ( + Operation::RemoveContainer, + ResourceTarget::Deployment(data.deployment.clone()), + ), + + // Build + ExecuteRequest::RunBuild(data) => ( + Operation::RunBuild, + ResourceTarget::Build(data.build.clone()), + ), + ExecuteRequest::CancelBuild(data) => ( + Operation::CancelBuild, + ResourceTarget::Build(data.build.clone()), + ), + + // Repo + ExecuteRequest::CloneRepo(data) => ( + Operation::CloneRepo, + ResourceTarget::Repo(data.repo.clone()), + ), + ExecuteRequest::PullRepo(data) => { + (Operation::PullRepo, ResourceTarget::Repo(data.repo.clone())) + } + + // Procedure + ExecuteRequest::RunProcedure(data) => ( + Operation::RunProcedure, + ResourceTarget::Procedure(data.procedure.clone()), + ), + + // Server template + ExecuteRequest::LaunchServer(data) => ( + Operation::LaunchServer, + ResourceTarget::ServerTemplate(data.server_template.clone()), + ), + }; + let mut update = make_update(target, operation, user); + update.in_progress(); + update.id = add_update(update.clone()).await?; + Ok(update) +} diff --git a/bin/core/src/listener/github.rs b/bin/core/src/listener/github.rs index 36e20476e..b26c40d17 100644 --- a/bin/core/src/listener/github.rs +++ b/bin/core/src/listener/github.rs @@ -18,7 +18,9 @@ use tracing::Instrument; use crate::{ config::core_config, - helpers::{cache::Cache, random_duration}, + helpers::{ + cache::Cache, random_duration, update::init_execution_update, + }, resource, state::State, }; @@ -137,12 +139,15 @@ async fn handle_build_webhook( if request_branch != build.config.branch { return Err(anyhow!("request branch does not match expected")); } - State - .resolve( - execute::RunBuild { build: build_id }, - github_user().to_owned(), - ) - .await?; + let user = github_user().to_owned(); + let req = crate::api::execute::ExecuteRequest::RunBuild( + execute::RunBuild { build: build_id }, + ); + let update = init_execution_update(&req, &user).await?; + let crate::api::execute::ExecuteRequest::RunBuild(req) = req else { + unreachable!() + }; + State.resolve(req, (user, update)).await?; Ok(()) } @@ -166,12 +171,16 @@ async fn handle_repo_clone_webhook( if request_branch != repo.config.branch { return Err(anyhow!("request branch does not match expected")); } - State - .resolve( - execute::CloneRepo { repo: repo_id }, - github_user().to_owned(), - ) - .await?; + let user = github_user().to_owned(); + let req = crate::api::execute::ExecuteRequest::CloneRepo( + execute::CloneRepo { repo: repo_id }, + ); + let update = init_execution_update(&req, &user).await?; + let crate::api::execute::ExecuteRequest::CloneRepo(req) = req + else { + unreachable!() + }; + State.resolve(req, (user, update)).await?; Ok(()) } @@ -195,12 +204,15 @@ async fn handle_repo_pull_webhook( if request_branch != repo.config.branch { return Err(anyhow!("request branch does not match expected")); } - State - .resolve( - execute::PullRepo { repo: repo_id }, - github_user().to_owned(), - ) - .await?; + let user = github_user().to_owned(); + let req = crate::api::execute::ExecuteRequest::PullRepo( + execute::PullRepo { repo: repo_id }, + ); + let update = init_execution_update(&req, &user).await?; + let crate::api::execute::ExecuteRequest::PullRepo(req) = req else { + unreachable!() + }; + State.resolve(req, (user, update)).await?; Ok(()) } @@ -226,14 +238,18 @@ async fn handle_procedure_webhook( if !procedure.config.webhook_enabled { return Err(anyhow!("procedure does not have webhook enabled")); } - State - .resolve( - execute::RunProcedure { - procedure: procedure_id, - }, - github_user().to_owned(), - ) - .await?; + let user = github_user().to_owned(); + let req = crate::api::execute::ExecuteRequest::RunProcedure( + execute::RunProcedure { + procedure: procedure_id, + }, + ); + let update = init_execution_update(&req, &user).await?; + let crate::api::execute::ExecuteRequest::RunProcedure(req) = req + else { + unreachable!() + }; + State.resolve(req, (user, update)).await?; Ok(()) } diff --git a/bin/migrator/src/legacy/v0/mod.rs b/bin/migrator/src/legacy/v0/mod.rs index 9ee518052..720631ac0 100644 --- a/bin/migrator/src/legacy/v0/mod.rs +++ b/bin/migrator/src/legacy/v0/mod.rs @@ -183,9 +183,9 @@ impl From for monitor_client::entities::Operation { Operation::CreateServer => CreateServer, Operation::UpdateServer => UpdateServer, Operation::DeleteServer => DeleteServer, - Operation::PruneImagesServer => PruneImagesServer, - Operation::PruneContainersServer => PruneContainersServer, - Operation::PruneNetworksServer => PruneNetworksServer, + Operation::PruneImagesServer => PruneImages, + Operation::PruneContainersServer => PruneContainers, + Operation::PruneNetworksServer => PruneNetworks, Operation::RenameServer => RenameServer, Operation::CreateBuild => CreateBuild, Operation::UpdateBuild => UpdateBuild, @@ -194,7 +194,7 @@ impl From for monitor_client::entities::Operation { Operation::CreateDeployment => CreateDeployment, Operation::UpdateDeployment => UpdateDeployment, Operation::DeleteDeployment => DeleteDeployment, - Operation::DeployContainer => DeployContainer, + Operation::DeployContainer => Deploy, Operation::StopContainer => StopContainer, Operation::StartContainer => StartContainer, Operation::RemoveContainer => RemoveContainer, diff --git a/client/core/rs/src/entities/mod.rs b/client/core/rs/src/entities/mod.rs index 7fdb472e1..fc9ab12ec 100644 --- a/client/core/rs/src/entities/mod.rs +++ b/client/core/rs/src/entities/mod.rs @@ -392,9 +392,9 @@ pub enum Operation { UpdateServer, DeleteServer, RenameServer, - PruneImagesServer, - PruneContainersServer, - PruneNetworksServer, + PruneImages, + PruneContainers, + PruneNetworks, CreateNetwork, DeleteNetwork, StopAllContainers, @@ -415,7 +415,7 @@ pub enum Operation { CreateDeployment, UpdateDeployment, DeleteDeployment, - DeployContainer, + Deploy, StopContainer, StartContainer, RemoveContainer,