diff --git a/bin/core/src/api/execute/procedure.rs b/bin/core/src/api/execute/procedure.rs index 2901ae0bb..7f4e75035 100644 --- a/bin/core/src/api/execute/procedure.rs +++ b/bin/core/src/api/execute/procedure.rs @@ -6,6 +6,7 @@ use monitor_client::{ update::Update, user::User, Operation, }, }; +use mungos::{by_id::update_one_by_id, mongodb::bson::to_document}; use resolver_api::Resolve; use serror::serialize_error_pretty; use tokio::sync::Mutex; @@ -14,7 +15,9 @@ use crate::{ helpers::{ procedure::execute_procedure, update::{add_update, make_update, update_update}, - }, resource, state::{action_states, State} + }, + resource::{self, refresh_procedure_state_cache}, + state::{action_states, db_client, State}, }; #[async_trait] @@ -74,6 +77,21 @@ impl Resolve for State { update.finalize(); + // Need to manually update the update before cache refresh, + // and before broadcast with add_update. + // The Err case of to_document should be unreachable, + // but will fail to update cache in that case. + if let Ok(update_doc) = to_document(&update) { + let _ = update_one_by_id( + &db_client().await.updates, + &update.id, + mungos::update::Update::Set(update_doc), + None, + ) + .await; + refresh_procedure_state_cache().await; + } + update_update(update.clone()).await?; Ok(update) diff --git a/bin/core/src/main.rs b/bin/core/src/main.rs index e3a543a2f..3c22216ff 100644 --- a/bin/core/src/main.rs +++ b/bin/core/src/main.rs @@ -37,6 +37,7 @@ async fn app() -> anyhow::Result<()> { helpers::prune::spawn_prune_loop(); resource::spawn_build_state_refresh_loop(); resource::spawn_repo_state_refresh_loop(); + resource::spawn_procedure_state_refresh_loop(); // Setup static frontend services let frontend_path = frontend_path(); diff --git a/bin/core/src/resource/mod.rs b/bin/core/src/resource/mod.rs index f6748337e..5816159ba 100644 --- a/bin/core/src/resource/mod.rs +++ b/bin/core/src/resource/mod.rs @@ -51,6 +51,9 @@ mod server_template; pub use build::{ refresh_build_state_cache, spawn_build_state_refresh_loop, }; +pub use procedure::{ + refresh_procedure_state_cache, spawn_procedure_state_refresh_loop, +}; pub use repo::{ refresh_repo_state_cache, spawn_repo_state_refresh_loop, }; diff --git a/bin/core/src/resource/procedure.rs b/bin/core/src/resource/procedure.rs index b45cfe269..ee3ef2980 100644 --- a/bin/core/src/resource/procedure.rs +++ b/bin/core/src/resource/procedure.rs @@ -1,4 +1,6 @@ -use anyhow::anyhow; +use std::time::Duration; + +use anyhow::{anyhow, Context}; use monitor_client::{ api::execute::Execution, entities::{ @@ -8,7 +10,7 @@ use monitor_client::{ procedure::{ PartialProcedureConfig, Procedure, ProcedureConfig, ProcedureConfigDiff, ProcedureListItem, ProcedureListItemInfo, - ProcedureQuerySpecifics, + ProcedureQuerySpecifics, ProcedureState, }, repo::Repo, resource::Resource, @@ -18,9 +20,12 @@ use monitor_client::{ Operation, }, }; -use mungos::mongodb::Collection; +use mungos::{ + find::find_collect, + mongodb::{bson::doc, options::FindOneOptions, Collection}, +}; -use crate::state::{action_states, db_client}; +use crate::state::{action_states, db_client, procedure_state_cache}; impl super::MonitorResource for Procedure { type Config = ProcedureConfig; @@ -42,6 +47,7 @@ impl super::MonitorResource for Procedure { async fn to_list_item( procedure: Resource, ) -> anyhow::Result { + let state = get_procedure_state(&procedure.id).await; Ok(ProcedureListItem { name: procedure.name, id: procedure.id, @@ -49,6 +55,7 @@ impl super::MonitorResource for Procedure { resource_type: ResourceTargetVariant::Procedure, info: ProcedureListItemInfo { procedure_type: procedure.config.procedure_type, + state, }, }) } @@ -260,3 +267,80 @@ async fn validate_config( } Ok(()) } + +pub fn spawn_procedure_state_refresh_loop() { + tokio::spawn(async move { + loop { + refresh_procedure_state_cache().await; + tokio::time::sleep(Duration::from_secs(60)).await; + } + }); +} + +pub async fn refresh_procedure_state_cache() { + let _ = async { + let procedures = + find_collect(&db_client().await.procedures, None, None) + .await + .context("failed to get procedures from db")?; + let cache = procedure_state_cache(); + for procedure in procedures { + let state = get_procedure_state_from_db(&procedure.id).await; + cache.insert(procedure.id, state).await; + } + anyhow::Ok(()) + } + .await + .inspect_err(|e| { + error!("failed to refresh build state cache | {e:#}") + }); +} + +async fn get_procedure_state(id: &String) -> ProcedureState { + if action_states() + .procedure + .get(id) + .await + .map(|s| s.get().map(|s| s.running)) + .transpose() + .ok() + .flatten() + .unwrap_or_default() + { + return ProcedureState::Running; + } + procedure_state_cache().get(id).await.unwrap_or_default() +} + +async fn get_procedure_state_from_db(id: &str) -> ProcedureState { + async { + let state = db_client() + .await + .updates + .find_one( + doc! { + "target.type": "Procedure", + "target.id": id, + "operation": "RunProcedure" + }, + FindOneOptions::builder() + .sort(doc! { "start_ts": -1 }) + .build(), + ) + .await? + .map(|u| { + if u.success { + ProcedureState::Ok + } else { + ProcedureState::Failed + } + }) + .unwrap_or(ProcedureState::Ok); + anyhow::Ok(state) + } + .await + .inspect_err(|e| { + warn!("failed to get procedure state for {id} | {e:#}") + }) + .unwrap_or(ProcedureState::Unknown) +} diff --git a/bin/core/src/state.rs b/bin/core/src/state.rs index 769e27c49..4e6eae358 100644 --- a/bin/core/src/state.rs +++ b/bin/core/src/state.rs @@ -1,7 +1,8 @@ use std::sync::{Arc, OnceLock}; use monitor_client::entities::{ - build::BuildState, deployment::DeploymentState, repo::RepoState, + build::BuildState, deployment::DeploymentState, + procedure::ProcedureState, repo::RepoState, }; use tokio::sync::OnceCell; @@ -80,3 +81,11 @@ pub fn repo_state_cache() -> &'static RepoStateCache { static REPO_STATE_CACHE: OnceLock = OnceLock::new(); REPO_STATE_CACHE.get_or_init(Default::default) } + +pub type ProcedureStateCache = Cache; + +pub fn procedure_state_cache() -> &'static ProcedureStateCache { + static PROCEDURE_STATE_CACHE: OnceLock = + OnceLock::new(); + PROCEDURE_STATE_CACHE.get_or_init(Default::default) +} diff --git a/client/core/rs/src/entities/procedure.rs b/client/core/rs/src/entities/procedure.rs index 7c5375c7f..afe33dc5b 100644 --- a/client/core/rs/src/entities/procedure.rs +++ b/client/core/rs/src/entities/procedure.rs @@ -20,7 +20,26 @@ pub type ProcedureListItem = ResourceListItem; #[typeshare] #[derive(Serialize, Deserialize, Debug, Clone)] pub struct ProcedureListItemInfo { + /// Sequence or Parallel. pub procedure_type: ProcedureType, + /// Reflect whether last run successful / currently running. + pub state: ProcedureState, +} + +#[typeshare] +#[derive( + Debug, Clone, Copy, Default, Serialize, Deserialize, Display, +)] +pub enum ProcedureState { + /// Last run successful + Ok, + /// Last run failed + Failed, + /// Currently running + Running, + /// Other case (never run) + #[default] + Unknown, } #[typeshare]