Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
125 changes: 97 additions & 28 deletions src/handlers/http/alerts.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,14 +22,18 @@ use crate::{
alerts::{
ALERTS, AlertError, AlertState, Severity,
alert_enums::{AlertType, NotificationState},
alert_structs::{AlertConfig, AlertRequest, AlertStateEntry, NotificationStateRequest},
alert_structs::{
AlertConfig, AlertConfigResponse, AlertRequest, AlertStateEntry,
NotificationStateRequest,
},
alert_traits::AlertTrait,
alert_types::ThresholdAlert,
target::Retry,
user_auth_for_alert_config,
},
metastore::metastore_traits::MetastoreObject,
parseable::PARSEABLE,
rbac::map::SessionKey,
utils::{actix::extract_session_key_from_req, get_tenant_id_from_request},
};
use actix_web::{
Expand Down Expand Up @@ -211,9 +215,16 @@ pub async fn list(req: HttpRequest) -> Result<impl Responder, AlertError> {
let session_key = extract_session_key_from_req(&req)?;
let query_map = web::Query::<HashMap<String, String>>::from_query(req.query_string())
.map_err(|_| AlertError::InvalidQueryParameter("malformed query parameters".to_string()))?;
Ok(web::Json(list_internal(session_key, &query_map).await?))
}

/// Shared alert-list implementation for HTTP handlers and in-process callers.
pub async fn list_internal(
session_key: SessionKey,
query_map: &HashMap<String, String>,
) -> Result<Vec<serde_json::Map<String, serde_json::Value>>, AlertError> {
// Parse and validate query parameters
let params = parse_list_query_params(&query_map)?;
let params = parse_list_query_params(query_map)?;

// Get alerts from the manager
let guard = ALERTS.read().await;
Expand Down Expand Up @@ -241,7 +252,7 @@ pub async fn list(req: HttpRequest) -> Result<impl Responder, AlertError> {
// Paginate results
let paginated_alerts = paginate_alerts(alerts_summary, params.offset, params.limit);

Ok(web::Json(paginated_alerts))
Ok(paginated_alerts)
}

// POST /alerts
Expand All @@ -250,6 +261,18 @@ pub async fn post(
Json(alert): Json<AlertRequest>,
) -> Result<impl Responder, AlertError> {
let tenant_id = get_tenant_id_from_request(&req);
let session_key = extract_session_key_from_req(&req)?;
Ok(web::Json(
post_internal(alert, &session_key, &tenant_id).await?,
))
}

/// Shared alert-creation implementation for HTTP handlers and in-process callers.
pub async fn post_internal(
alert: AlertRequest,
session_key: &SessionKey,
tenant_id: &Option<String>,
) -> Result<AlertConfigResponse, AlertError> {
let mut alert: AlertConfig = alert.into(tenant_id.clone()).await?;

if alert.notification_config.interval > alert.get_eval_frequency() {
Expand Down Expand Up @@ -303,22 +326,20 @@ pub async fn post(

// validate the incoming alert query
// does the user have access to these tables or not?
let session_key = extract_session_key_from_req(&req)?;

alert.validate(&session_key).await?;
alert.validate(session_key).await?;

// update persistent storage first
PARSEABLE
.metastore
.put_alert(&alert.to_alert_config(), &tenant_id)
.put_alert(&alert.to_alert_config(), tenant_id)
.await?;

// create initial alert state entry (default to NotTriggered)
let state_entry =
AlertStateEntry::new(*alert.get_id(), AlertState::NotTriggered, tenant_id.clone());
PARSEABLE
.metastore
.put_alert_state(&state_entry as &dyn MetastoreObject, &tenant_id)
.put_alert_state(&state_entry as &dyn MetastoreObject, tenant_id)
.await?;

// update in memory
Expand All @@ -327,26 +348,37 @@ pub async fn post(
// start the task
alerts.start_task(alert.clone_box()).await?;

Ok(web::Json(alert.to_alert_config().to_response()))
Ok(alert.to_alert_config().to_response())
}

// GET /alerts/{alert_id}
pub async fn get(req: HttpRequest, alert_id: Path<Ulid>) -> Result<impl Responder, AlertError> {
let session_key = extract_session_key_from_req(&req)?;
let alert_id = alert_id.into_inner();
let tenant_id = get_tenant_id_from_request(&req);
Ok(web::Json(
get_internal(&session_key, alert_id, &tenant_id).await?,
))
}

/// Shared alert lookup implementation for HTTP handlers and in-process callers.
pub async fn get_internal(
session_key: &SessionKey,
alert_id: Ulid,
tenant_id: &Option<String>,
) -> Result<AlertConfigResponse, AlertError> {
let guard = ALERTS.read().await;
let alerts = if let Some(alerts) = guard.as_ref() {
alerts
} else {
return Err(AlertError::CustomError("No AlertManager set".into()));
};

let alert = alerts.get_alert_by_id(alert_id, &tenant_id).await?;
let alert = alerts.get_alert_by_id(alert_id, tenant_id).await?;
// validate that the user has access to the tables mentioned in the query
user_auth_for_alert_config(&session_key, &alert.to_alert_config()).await?;
user_auth_for_alert_config(session_key, &alert.to_alert_config()).await?;

Ok(web::Json(alert.to_alert_config().to_response()))
Ok(alert.to_alert_config().to_response())
}

// DELETE /alerts/{alert_id}
Expand Down Expand Up @@ -457,6 +489,17 @@ pub async fn disable_alert(
let session_key = extract_session_key_from_req(&req)?;
let alert_id = alert_id.into_inner();
let tenant_id = get_tenant_id_from_request(&req);
Ok(web::Json(
disable_alert_internal(&session_key, alert_id, &tenant_id).await?,
))
}

/// Shared alert-disable implementation for HTTP handlers and in-process callers.
pub async fn disable_alert_internal(
session_key: &SessionKey,
alert_id: Ulid,
tenant_id: &Option<String>,
) -> Result<AlertConfigResponse, AlertError> {
let guard = ALERTS.write().await;
let alerts = if let Some(alerts) = guard.as_ref() {
alerts
Expand All @@ -465,16 +508,16 @@ pub async fn disable_alert(
};

// check if alert id exists in map
let alert = alerts.get_alert_by_id(alert_id, &tenant_id).await?;
let alert = alerts.get_alert_by_id(alert_id, tenant_id).await?;
// validate that the user has access to the tables mentioned in the query
user_auth_for_alert_config(&session_key, &alert.to_alert_config()).await?;
user_auth_for_alert_config(session_key, &alert.to_alert_config()).await?;

alerts
.update_state(alert_id, AlertState::Disabled, Some("".into()), &tenant_id)
.update_state(alert_id, AlertState::Disabled, Some("".into()), tenant_id)
.await?;
let alert = alerts.get_alert_by_id(alert_id, &tenant_id).await?;
let alert = alerts.get_alert_by_id(alert_id, tenant_id).await?;

Ok(web::Json(alert.to_alert_config().to_response()))
Ok(alert.to_alert_config().to_response())
}

// PATCH /alerts/{alert_id}/enable
Expand All @@ -487,6 +530,17 @@ pub async fn enable_alert(
let session_key = extract_session_key_from_req(&req)?;
let alert_id = alert_id.into_inner();
let tenant_id = get_tenant_id_from_request(&req);
Ok(web::Json(
enable_alert_internal(&session_key, alert_id, &tenant_id).await?,
))
}

/// Shared alert-enable implementation for HTTP handlers and in-process callers.
pub async fn enable_alert_internal(
session_key: &SessionKey,
alert_id: Ulid,
tenant_id: &Option<String>,
) -> Result<AlertConfigResponse, AlertError> {
let guard = ALERTS.write().await;
let alerts = if let Some(alerts) = guard.as_ref() {
alerts
Expand All @@ -495,7 +549,7 @@ pub async fn enable_alert(
};

// check if alert id exists in map
let alert = alerts.get_alert_by_id(alert_id, &tenant_id).await?;
let alert = alerts.get_alert_by_id(alert_id, tenant_id).await?;

// only run if alert is disabled
if alert.get_state().ne(&AlertState::Disabled) {
Expand All @@ -505,19 +559,19 @@ pub async fn enable_alert(
}

// validate that the user has access to the tables mentioned in the query
user_auth_for_alert_config(&session_key, &alert.to_alert_config()).await?;
user_auth_for_alert_config(session_key, &alert.to_alert_config()).await?;

alerts
.update_state(
alert_id,
AlertState::NotTriggered,
Some("".into()),
&tenant_id,
tenant_id,
)
.await?;
let alert = alerts.get_alert_by_id(alert_id, &tenant_id).await?;
let alert = alerts.get_alert_by_id(alert_id, tenant_id).await?;

Ok(web::Json(alert.to_alert_config().to_response()))
Ok(alert.to_alert_config().to_response())
}

// PUT /alerts/{alert_id}
Expand Down Expand Up @@ -616,16 +670,27 @@ pub async fn evaluate_alert(
let session_key = extract_session_key_from_req(&req)?;
let alert_id = alert_id.into_inner();
let tenant_id = get_tenant_id_from_request(&req);
Ok(Json(
evaluate_alert_internal(&session_key, alert_id, &tenant_id).await?,
))
}

/// Shared alert-evaluation implementation for HTTP handlers and in-process callers.
pub async fn evaluate_alert_internal(
session_key: &SessionKey,
alert_id: Ulid,
tenant_id: &Option<String>,
) -> Result<AlertConfigResponse, AlertError> {
let guard = ALERTS.write().await;
let alerts = if let Some(alerts) = guard.as_ref() {
alerts
} else {
return Err(AlertError::CustomError("No AlertManager set".into()));
};

let alert = alerts.get_alert_by_id(alert_id, &tenant_id).await?;
let alert = alerts.get_alert_by_id(alert_id, tenant_id).await?;

user_auth_for_alert_config(&session_key, &alert.to_alert_config()).await?;
user_auth_for_alert_config(session_key, &alert.to_alert_config()).await?;

let config = alert.to_alert_config().to_response();

Expand All @@ -635,17 +700,21 @@ pub async fn evaluate_alert(
// add the task back again so that it evaluates right now
alerts.start_task(alert).await?;

Ok(Json(config))
Ok(config)
}

pub async fn list_tags(req: HttpRequest) -> Result<impl Responder, AlertError> {
let tenant_id = get_tenant_id_from_request(&req);
Ok(web::Json(list_tags_internal(&tenant_id).await?))
}

/// Shared alert-tag implementation for HTTP handlers and in-process callers.
pub async fn list_tags_internal(tenant_id: &Option<String>) -> Result<Vec<String>, AlertError> {
let guard = ALERTS.read().await;
let alerts = if let Some(alerts) = guard.as_ref() {
alerts
} else {
return Err(AlertError::CustomError("No AlertManager set".into()));
};
let tenant_id = get_tenant_id_from_request(&req);
let tags = alerts.list_tags(&tenant_id).await;
Ok(web::Json(tags))
Ok(alerts.list_tags(tenant_id).await)
}
26 changes: 19 additions & 7 deletions src/handlers/http/cluster/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -939,7 +939,14 @@ pub async fn send_retention_cleanup_request(

/// Fetches cluster information for all nodes (ingestor, indexer, querier and prism)
pub async fn get_cluster_info(req: HttpRequest) -> Result<impl Responder, StreamError> {
let tenant_id = &get_tenant_id_from_request(&req);
let tenant_id = get_tenant_id_from_request(&req);
Ok(actix_web::HttpResponse::Ok().json(get_cluster_info_internal(&tenant_id).await?))
}

/// Shared cluster-info implementation for HTTP handlers and in-process callers.
pub async fn get_cluster_info_internal(
tenant_id: &Option<String>,
) -> Result<Vec<utils::ClusterInfo>, StreamError> {
// Get querier, ingestor and indexer metadata concurrently
let (prism_result, querier_result, ingestor_result, indexer_result) = future::join4(
get_node_info(NodeType::Prism, tenant_id),
Expand Down Expand Up @@ -996,7 +1003,7 @@ pub async fn get_cluster_info(req: HttpRequest) -> Result<impl Responder, Stream
infos.extend(querier_infos?);
infos.extend(ingestor_infos?);
infos.extend(indexer_infos?);
Ok(actix_web::HttpResponse::Ok().json(infos))
Ok(infos)
}

/// Fetches info for a single node
Expand Down Expand Up @@ -1085,13 +1092,18 @@ async fn fetch_nodes_info<T: Metadata>(
}

pub async fn get_cluster_metrics(req: HttpRequest) -> Result<impl Responder, PostError> {
let tenant_id = &get_tenant_id_from_request(&req);
let dresses = fetch_cluster_metrics(tenant_id).await.map_err(|err| {
let tenant_id = get_tenant_id_from_request(&req);
Ok(actix_web::HttpResponse::Ok().json(get_cluster_metrics_internal(&tenant_id).await?))
}

/// Shared cluster-metrics implementation for HTTP handlers and in-process callers.
pub async fn get_cluster_metrics_internal(
tenant_id: &Option<String>,
) -> Result<Vec<Metrics>, PostError> {
fetch_cluster_metrics(tenant_id).await.map_err(|err| {
error!("Fatal: failed to fetch cluster metrics: {:?}", err);
PostError::Invalid(err.into())
})?;

Ok(actix_web::HttpResponse::Ok().json(dresses))
})
}

/// get node info for a specific node type
Expand Down
16 changes: 12 additions & 4 deletions src/handlers/http/health_check.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,11 @@ use crate::{parseable::PARSEABLE, storage::object_storage::sync_all_streams};
pub static SIGNAL_RECEIVED: Lazy<Arc<Mutex<bool>>> = Lazy::new(|| Arc::new(Mutex::new(false)));

pub async fn liveness() -> HttpResponse {
HttpResponse::new(StatusCode::OK)
HttpResponse::new(liveness_internal())
}

pub fn liveness_internal() -> StatusCode {
StatusCode::OK
}

pub async fn check_shutdown_middleware(
Expand Down Expand Up @@ -107,16 +111,20 @@ async fn perform_object_store_sync() {

pub async fn readiness(req: HttpRequest) -> HttpResponse {
let tenant_id = get_tenant_id_from_request(&req);
HttpResponse::new(readiness_internal(&tenant_id).await)
}

pub async fn readiness_internal(tenant_id: &Option<String>) -> StatusCode {
// Check the object store connection
if PARSEABLE
.storage
.get_object_store()
.check(&tenant_id)
.check(tenant_id)
.await
.is_ok()
{
HttpResponse::new(StatusCode::OK)
StatusCode::OK
} else {
HttpResponse::new(StatusCode::SERVICE_UNAVAILABLE)
StatusCode::SERVICE_UNAVAILABLE
}
}
Loading
Loading