mirror of
https://github.com/moghtech/komodo.git
synced 2026-09-07 16:01:24 +00:00
add ProcedureState
This commit is contained in:
@@ -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<RunProcedure, User> 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)
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -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,
|
||||
};
|
||||
|
||||
@@ -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<Self::Config, Self::Info>,
|
||||
) -> anyhow::Result<Self::ListItem> {
|
||||
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)
|
||||
}
|
||||
|
||||
+10
-1
@@ -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<RepoStateCache> = OnceLock::new();
|
||||
REPO_STATE_CACHE.get_or_init(Default::default)
|
||||
}
|
||||
|
||||
pub type ProcedureStateCache = Cache<String, ProcedureState>;
|
||||
|
||||
pub fn procedure_state_cache() -> &'static ProcedureStateCache {
|
||||
static PROCEDURE_STATE_CACHE: OnceLock<ProcedureStateCache> =
|
||||
OnceLock::new();
|
||||
PROCEDURE_STATE_CACHE.get_or_init(Default::default)
|
||||
}
|
||||
|
||||
@@ -20,7 +20,26 @@ pub type ProcedureListItem = ResourceListItem<ProcedureListItemInfo>;
|
||||
#[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]
|
||||
|
||||
Reference in New Issue
Block a user