execute api returns update immediately

This commit is contained in:
mbecker20
2024-06-02 04:14:51 -07:00
parent f1e51d275c
commit f5ce3570e4
14 changed files with 398 additions and 290 deletions
Generated
+6 -6
View File
@@ -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",
+1 -1
View File
@@ -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"
+28 -20
View File
@@ -52,12 +52,14 @@ use crate::{
state::{action_states, db_client, State},
};
impl Resolve<RunBuild, User> for State {
use crate::helpers::update::init_execution_update;
impl Resolve<RunBuild, (User, Update)> for State {
#[instrument(name = "RunBuild", skip(self, user))]
async fn resolve(
&self,
RunBuild { build }: RunBuild,
user: User,
(user, mut update): (User, Update),
) -> anyhow::Result<Update> {
let mut build = resource::get_check_permissions::<Build>(
&build,
@@ -77,8 +79,6 @@ impl Resolve<RunBuild, User> 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<CancelBuild, User> for State {
impl Resolve<CancelBuild, (User, Update)> for State {
#[instrument(name = "CancelBuild", skip(self, user))]
async fn resolve(
&self,
CancelBuild { build }: CancelBuild,
user: User,
(user, mut update): (User, Update),
) -> anyhow::Result<CancelBuildResponse> {
let build = resource::get_check_permissions::<Build>(
&build,
@@ -370,9 +370,6 @@ impl Resolve<CancelBuild, User> 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))),
}
}
}
+39 -70
View File
@@ -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<Deploy, User> for State {
use crate::helpers::update::init_execution_update;
impl Resolve<Deploy, (User, Update)> for State {
#[instrument(name = "Deploy", skip(self, user))]
async fn resolve(
&self,
@@ -41,7 +43,7 @@ impl Resolve<Deploy, User> for State {
stop_signal,
stop_time,
}: Deploy,
user: User,
(user, mut update): (User, Update),
) -> anyhow::Result<Update> {
let mut deployment =
resource::get_check_permissions::<Deployment>(
@@ -129,11 +131,6 @@ impl Resolve<Deploy, User> 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<Deploy, User> 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<Deploy, User> for State {
}
}
impl Resolve<StartContainer, User> for State {
impl Resolve<StartContainer, (User, Update)> for State {
#[instrument(name = "StartContainer", skip(self, user))]
async fn resolve(
&self,
StartContainer { deployment }: StartContainer,
user: User,
(user, mut update): (User, Update),
) -> anyhow::Result<Update> {
let deployment = resource::get_check_permissions::<Deployment>(
&deployment,
@@ -230,20 +228,6 @@ impl Resolve<StartContainer, User> 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<StartContainer, User> 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<StopContainer, User> for State {
impl Resolve<StopContainer, (User, Update)> for State {
#[instrument(name = "StopContainer", skip(self, user))]
async fn resolve(
&self,
@@ -274,7 +258,7 @@ impl Resolve<StopContainer, User> for State {
signal,
time,
}: StopContainer,
user: User,
(user, mut update): (User, Update),
) -> anyhow::Result<Update> {
let deployment = resource::get_check_permissions::<Deployment>(
&deployment,
@@ -308,11 +292,6 @@ impl Resolve<StopContainer, User> 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<StopContainer, User> 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<StopAllContainers, User> for State {
impl Resolve<StopAllContainers, (User, Update)> for State {
#[instrument(name = "StopAllContainers", skip(self, user))]
async fn resolve(
&self,
StopAllContainers { server }: StopAllContainers,
user: User,
(user, mut update): (User, Update),
) -> anyhow::Result<Update> {
let (server, status) = get_server_with_status(&server).await?;
if status != ServerState::Ok {
@@ -375,23 +354,27 @@ impl Resolve<StopAllContainers, User> 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<StopAllContainers, User> for State {
}
}
impl Resolve<RemoveContainer, User> for State {
impl Resolve<RemoveContainer, (User, Update)> for State {
#[instrument(name = "RemoveContainer", skip(self, user))]
async fn resolve(
&self,
@@ -431,7 +414,7 @@ impl Resolve<RemoveContainer, User> for State {
signal,
time,
}: RemoveContainer,
user: User,
(user, mut update): (User, Update),
) -> anyhow::Result<Update> {
let deployment = resource::get_check_permissions::<Deployment>(
&deployment,
@@ -465,20 +448,6 @@ impl Resolve<RemoveContainer, User> 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(),
+50 -25
View File
@@ -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<User>,
Json(request): Json<ExecuteRequest>,
) -> serror::Result<(TypedHeader<ContentType>, String)> {
) -> serror::Result<Json<Update>> {
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<String> {
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:#}");
+6 -18
View File
@@ -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<RunProcedure, User> for State {
impl Resolve<RunProcedure, (User, Update)> for State {
#[instrument(name = "RunProcedure", skip(self, user))]
async fn resolve(
&self,
RunProcedure { procedure }: RunProcedure,
user: User,
(user, update): (User, Update),
) -> anyhow::Result<Update> {
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<Output = anyhow::Result<Update>> + 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;
+6 -38
View File
@@ -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<CloneRepo, User> for State {
impl Resolve<CloneRepo, (User, Update)> for State {
#[instrument(name = "CloneRepo", skip(self, user))]
async fn resolve(
&self,
CloneRepo { repo }: CloneRepo,
user: User,
(user, mut update): (User, Update),
) -> anyhow::Result<Update> {
let repo = resource::get_check_permissions::<Repo>(
&repo,
@@ -61,20 +57,6 @@ impl Resolve<CloneRepo, User> 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<CloneRepo, User> for State {
}
}
impl Resolve<PullRepo, User> for State {
impl Resolve<PullRepo, (User, Update)> for State {
#[instrument(name = "PullRepo", skip(self, user))]
async fn resolve(
&self,
PullRepo { repo }: PullRepo,
user: User,
(user, mut update): (User, Update),
) -> anyhow::Result<Update> {
let repo = resource::get_check_permissions::<Repo>(
&repo,
@@ -136,20 +118,6 @@ impl Resolve<PullRepo, User> 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(),
+7 -26
View File
@@ -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<PruneContainers, User> for State {
impl Resolve<PruneContainers, (User, Update)> for State {
#[instrument(name = "PruneContainers", skip(self, user))]
async fn resolve(
&self,
PruneContainers { server }: PruneContainers,
user: User,
(user, mut update): (User, Update),
) -> anyhow::Result<Update> {
let server = resource::get_check_permissions::<Server>(
&server,
@@ -50,11 +46,6 @@ impl Resolve<PruneContainers, User> 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<PruneContainers, User> for State {
}
}
impl Resolve<PruneNetworks, User> for State {
impl Resolve<PruneNetworks, (User, Update)> for State {
#[instrument(name = "PruneNetworks", skip(self, user))]
async fn resolve(
&self,
PruneNetworks { server }: PruneNetworks,
user: User,
(user, mut update): (User, Update),
) -> anyhow::Result<Update> {
let server = resource::get_check_permissions::<Server>(
&server,
@@ -106,11 +97,6 @@ impl Resolve<PruneNetworks, User> 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<PruneNetworks, User> for State {
}
}
impl Resolve<PruneImages, User> for State {
impl Resolve<PruneImages, (User, Update)> for State {
#[instrument(name = "PruneImages", skip(self, user))]
async fn resolve(
&self,
PruneImages { server }: PruneImages,
user: User,
(user, mut update): (User, Update),
) -> anyhow::Result<Update> {
let server = resource::get_check_permissions::<Server>(
&server,
@@ -162,11 +148,6 @@ impl Resolve<PruneImages, User> 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,
+4 -8
View File
@@ -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<LaunchServer, User> for State {
impl Resolve<LaunchServer, (User, Update)> for State {
#[instrument(name = "LaunchServer", skip(self, user))]
async fn resolve(
&self,
@@ -29,7 +28,7 @@ impl Resolve<LaunchServer, User> for State {
name,
server_template,
}: LaunchServer,
user: User,
(user, mut update): (User, Update),
) -> anyhow::Result<Update> {
// validate name isn't already taken by another server
if db_client()
@@ -55,14 +54,11 @@ impl Resolve<LaunchServer, User> 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) => {
+121 -42
View File
@@ -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(())
+79 -1
View File
@@ -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<Update> {
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)
}
+43 -27
View File
@@ -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(())
}
+4 -4
View File
@@ -183,9 +183,9 @@ impl From<Operation> 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<Operation> 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,
+4 -4
View File
@@ -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,