From 936ed6408a7eb600aa57e426286e2da379472586 Mon Sep 17 00:00:00 2001 From: Nikhil Sinha Date: Fri, 2 Oct 2026 10:19:25 +0530 Subject: [PATCH 1/5] abstract internal processing from handler function --- src/handlers/http/alerts.rs | 125 ++++++++++++++---- src/handlers/http/cluster/mod.rs | 26 +++- src/handlers/http/health_check.rs | 16 ++- src/handlers/http/logstream.rs | 50 ++++--- .../http/modal/query/querier_logstream.rs | 92 +++++++------ src/handlers/http/query.rs | 26 ++-- src/handlers/http/rbac.rs | 25 +++- src/handlers/http/role.rs | 29 +++- src/handlers/http/targets.rs | 37 +++++- src/handlers/http/traces.rs | 87 +++++++++--- src/handlers/http/users/dashboards.rs | 46 ++++--- src/utils/mod.rs | 10 +- 12 files changed, 409 insertions(+), 160 deletions(-) diff --git a/src/handlers/http/alerts.rs b/src/handlers/http/alerts.rs index 6fe8149b7..975eb479c 100644 --- a/src/handlers/http/alerts.rs +++ b/src/handlers/http/alerts.rs @@ -22,7 +22,10 @@ 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, @@ -30,6 +33,7 @@ use crate::{ }, metastore::metastore_traits::MetastoreObject, parseable::PARSEABLE, + rbac::map::SessionKey, utils::{actix::extract_session_key_from_req, get_tenant_id_from_request}, }; use actix_web::{ @@ -211,9 +215,16 @@ pub async fn list(req: HttpRequest) -> Result { let session_key = extract_session_key_from_req(&req)?; let query_map = web::Query::>::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, +) -> Result>, 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; @@ -241,7 +252,7 @@ pub async fn list(req: HttpRequest) -> Result { // Paginate results let paginated_alerts = paginate_alerts(alerts_summary, params.offset, params.limit); - Ok(web::Json(paginated_alerts)) + Ok(paginated_alerts) } // POST /alerts @@ -250,6 +261,18 @@ pub async fn post( Json(alert): Json, ) -> Result { 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, +) -> Result { let mut alert: AlertConfig = alert.into(tenant_id.clone()).await?; if alert.notification_config.interval > alert.get_eval_frequency() { @@ -303,14 +326,12 @@ 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) @@ -318,7 +339,7 @@ pub async fn post( 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 @@ -327,7 +348,7 @@ 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} @@ -335,6 +356,17 @@ pub async fn get(req: HttpRequest, alert_id: Path) -> Result, +) -> Result { let guard = ALERTS.read().await; let alerts = if let Some(alerts) = guard.as_ref() { alerts @@ -342,11 +374,11 @@ pub async fn get(req: HttpRequest, alert_id: Path) -> Result, +) -> Result { let guard = ALERTS.write().await; let alerts = if let Some(alerts) = guard.as_ref() { alerts @@ -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 @@ -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, +) -> Result { let guard = ALERTS.write().await; let alerts = if let Some(alerts) = guard.as_ref() { alerts @@ -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) { @@ -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} @@ -616,6 +670,17 @@ 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, +) -> Result { let guard = ALERTS.write().await; let alerts = if let Some(alerts) = guard.as_ref() { alerts @@ -623,9 +688,9 @@ pub async fn evaluate_alert( 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(); @@ -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 { + 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) -> Result, 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) } diff --git a/src/handlers/http/cluster/mod.rs b/src/handlers/http/cluster/mod.rs index ecae80cdc..8ecc5538f 100644 --- a/src/handlers/http/cluster/mod.rs +++ b/src/handlers/http/cluster/mod.rs @@ -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 { - 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, +) -> Result, 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), @@ -996,7 +1003,7 @@ pub async fn get_cluster_info(req: HttpRequest) -> Result( } pub async fn get_cluster_metrics(req: HttpRequest) -> Result { - 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, +) -> Result, 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 diff --git a/src/handlers/http/health_check.rs b/src/handlers/http/health_check.rs index bdae5097c..7ea5eea32 100644 --- a/src/handlers/http/health_check.rs +++ b/src/handlers/http/health_check.rs @@ -41,7 +41,11 @@ use crate::{parseable::PARSEABLE, storage::object_storage::sync_all_streams}; pub static SIGNAL_RECEIVED: Lazy>> = 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( @@ -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) -> 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 } } diff --git a/src/handlers/http/logstream.rs b/src/handlers/http/logstream.rs index ab6fd1a8b..6e427e8a3 100644 --- a/src/handlers/http/logstream.rs +++ b/src/handlers/http/logstream.rs @@ -177,23 +177,26 @@ pub async fn get_schema( ) -> Result { let stream_name = logstream.into_inner(); let tenant_id = get_tenant_id_from_request(&req); + let schema = get_schema_internal(&stream_name, &tenant_id).await?; + Ok((web::Json(schema), StatusCode::OK)) +} + +/// Shared schema implementation for HTTP handlers and in-process callers. +pub async fn get_schema_internal( + stream_name: &str, + tenant_id: &Option, +) -> Result, StreamError> { // Ensure parseable is aware of stream in distributed mode - if !PARSEABLE - .check_or_load_stream(&stream_name, &tenant_id) - .await - { - return Err(StreamNotFound(stream_name.clone()).into()); + if !PARSEABLE.check_or_load_stream(stream_name, tenant_id).await { + return Err(StreamNotFound(stream_name.to_owned()).into()); } - let stream = PARSEABLE.get_stream(&stream_name, &tenant_id)?; + let stream = PARSEABLE.get_stream(stream_name, tenant_id)?; if stream.is_deleting() { - return Err(StreamNotFound(stream_name.clone()).into()); + return Err(StreamNotFound(stream_name.to_owned()).into()); } - match update_schema_when_distributed(&vec![stream_name.clone()], &tenant_id).await { - Ok(_) => { - let schema = stream.get_schema(); - Ok((web::Json(schema), StatusCode::OK)) - } + match update_schema_when_distributed(&vec![stream_name.to_owned()], tenant_id).await { + Ok(_) => Ok(stream.get_schema()), Err(err) => Err(StreamError::Custom { msg: err.to_string(), status: StatusCode::EXPECTATION_FAILED, @@ -222,21 +225,26 @@ pub async fn get_retention( ) -> Result { let stream_name = stream_name.into_inner(); let tenant_id = get_tenant_id_from_request(&req); + let retention = get_retention_internal(&stream_name, &tenant_id).await?; + Ok((web::Json(retention), StatusCode::OK)) +} + +/// Shared retention implementation for HTTP handlers and in-process callers. +pub async fn get_retention_internal( + stream_name: &str, + tenant_id: &Option, +) -> Result { // For query mode, if the stream not found in memory map, // check if it exists in the storage // create stream and schema from storage - if !PARSEABLE - .check_or_load_stream(&stream_name, &tenant_id) - .await - { - return Err(StreamNotFound(stream_name.clone()).into()); + if !PARSEABLE.check_or_load_stream(stream_name, tenant_id).await { + return Err(StreamNotFound(stream_name.to_owned()).into()); } - let retention = PARSEABLE - .get_stream(&stream_name, &tenant_id)? + Ok(PARSEABLE + .get_stream(stream_name, tenant_id)? .get_retention() - .unwrap_or_default(); - Ok((web::Json(retention), StatusCode::OK)) + .unwrap_or_default()) } pub async fn put_retention( diff --git a/src/handlers/http/modal/query/querier_logstream.rs b/src/handlers/http/modal/query/querier_logstream.rs index 2bb170104..ec6365e97 100644 --- a/src/handlers/http/modal/query/querier_logstream.rs +++ b/src/handlers/http/modal/query/querier_logstream.rs @@ -174,60 +174,78 @@ pub async fn get_stats( ) -> Result { let stream_name = stream_name.into_inner(); let tenant_id = get_tenant_id_from_request(&req); + let query_map = web::Query::>::from_query(req.query_string()) + .map_err(|_| StreamError::InvalidQueryParameter(STATS_DATE_QUERY_PARAM.to_string()))?; + let date = if query_map.is_empty() { + None + } else { + Some( + query_map + .get(STATS_DATE_QUERY_PARAM) + .ok_or_else(|| { + StreamError::InvalidQueryParameter(STATS_DATE_QUERY_PARAM.to_string()) + })? + .as_str(), + ) + }; + Ok(web::Json( + get_stats_internal(&stream_name, &tenant_id, date).await?, + )) +} + +/// Shared statistics implementation for HTTP handlers and in-process callers. +pub async fn get_stats_internal( + stream_name: &str, + tenant_id: &Option, + date: Option<&str>, +) -> Result { // if the stream not found in memory map, //check if it exists in the storage //create stream and schema from storage - if !PARSEABLE.streams.contains(&stream_name, &tenant_id) + if !PARSEABLE.streams.contains(stream_name, tenant_id) && !PARSEABLE - .create_stream_and_schema_from_storage(&stream_name, &tenant_id) + .create_stream_and_schema_from_storage(stream_name, tenant_id) .await .unwrap_or(false) { - return Err(StreamNotFound(stream_name.clone()).into()); + return Err(StreamNotFound(stream_name.to_owned()).into()); } - let query_map = web::Query::>::from_query(req.query_string()) - .map_err(|_| StreamError::InvalidQueryParameter(STATS_DATE_QUERY_PARAM.to_string()))?; - - if !query_map.is_empty() { - let date_value = query_map.get(STATS_DATE_QUERY_PARAM).ok_or_else(|| { - StreamError::InvalidQueryParameter(STATS_DATE_QUERY_PARAM.to_string()) - })?; - - if !date_value.is_empty() { - let obs = PARSEABLE - .metastore - .get_all_stream_jsons(&stream_name, None, &tenant_id, false) - .await?; + if let Some(date_value) = date + && !date_value.is_empty() + { + let obs = PARSEABLE + .metastore + .get_all_stream_jsons(stream_name, None, tenant_id, false) + .await?; - let mut stream_jsons = Vec::new(); - for ob in obs { - let stream_metadata: ObjectStoreFormat = match serde_json::from_slice(&ob) { - Ok(d) => d, - Err(e) => { - error!("Failed to parse stream metadata: {:?}", e); - continue; - } - }; - stream_jsons.push(stream_metadata); - } + let mut stream_jsons = Vec::new(); + for ob in obs { + let stream_metadata: ObjectStoreFormat = match serde_json::from_slice(&ob) { + Ok(d) => d, + Err(e) => { + error!("Failed to parse stream metadata: {:?}", e); + continue; + } + }; + stream_jsons.push(stream_metadata); + } - let stats = fetch_daily_stats(date_value, &stream_jsons)?; + let stats = fetch_daily_stats(date_value, &stream_jsons)?; - let stats = serde_json::to_value(stats)?; + let stats = serde_json::to_value(stats)?; - return Ok(web::Json(stats)); - } + return Ok(stats); } - let stats = stats::get_current_stats(&stream_name, "json", &tenant_id) - .ok_or_else(|| StreamNotFound(stream_name.clone()))?; + let stats = stats::get_current_stats(stream_name, "json", tenant_id) + .ok_or_else(|| StreamNotFound(stream_name.to_owned()))?; let ingestor_stats = if PARSEABLE - .get_stream(&stream_name, &tenant_id) + .get_stream(stream_name, tenant_id) .is_ok_and(|stream| stream.get_stream_type() == StreamType::UserDefined) { - Some(fetch_stats_from_ingestors(&stream_name, &tenant_id).await?) + Some(fetch_stats_from_ingestors(stream_name, tenant_id).await?) } else { None }; @@ -251,7 +269,7 @@ pub async fn get_stats( "parquet", ); - QueriedStats::new(&stream_name, time, ingestion_stats, storage_stats) + QueriedStats::new(stream_name, time, ingestion_stats, storage_stats) }; let stats = if let Some(mut ingestor_stats) = ingestor_stats { @@ -264,5 +282,5 @@ pub async fn get_stats( let stats = serde_json::to_value(stats)?; - Ok(web::Json(stats)) + Ok(stats) } diff --git a/src/handlers/http/query.rs b/src/handlers/http/query.rs index bdfe3b6ef..1ab6f3aef 100644 --- a/src/handlers/http/query.rs +++ b/src/handlers/http/query.rs @@ -155,14 +155,23 @@ pub async fn get_records_and_fields_for_authorized_query( } pub async fn query(req: HttpRequest, query_request: Query) -> Result { + let tenant_id = get_tenant_id_from_request(&req); + let creds = extract_session_key_from_req(&req)?; + query_internal(query_request, &creds, &tenant_id).await +} + +/// Shared query implementation for HTTP handlers and in-process callers. +pub async fn query_internal( + query_request: Query, + creds: &SessionKey, + tenant_id: &Option, +) -> Result { let mut session_state = QUERY_SESSION.get_ctx().state(); let time_range = TimeRange::parse_human_time(&query_request.start_time, &query_request.end_time)?; let tables = resolve_stream_names(&query_request.query)?; // check or load streams in memory - create_streams_for_distributed(tables.clone(), &get_tenant_id_from_request(&req)).await?; - - let tenant_id = get_tenant_id_from_request(&req); + create_streams_for_distributed(tables.clone(), tenant_id).await?; session_state .config_mut() .options_mut() @@ -170,10 +179,9 @@ pub async fn query(req: HttpRequest, query_request: Query) -> Result Result`) diff --git a/src/handlers/http/rbac.rs b/src/handlers/http/rbac.rs index d729b1c46..c1dcab41d 100644 --- a/src/handlers/http/rbac.rs +++ b/src/handlers/http/rbac.rs @@ -72,7 +72,12 @@ impl From<&user::User> for User { // returns list of all registered users pub async fn list_users(req: HttpRequest) -> impl Responder { let tenant_id = get_tenant_id_from_request(&req); - web::Json(Users.collect_user::(&tenant_id)) + web::Json(list_users_internal(&tenant_id)) +} + +/// Shared user-list implementation for HTTP handlers and in-process callers. +pub fn list_users_internal(tenant_id: &Option) -> serde_json::Value { + serde_json::to_value(Users.collect_user::(tenant_id)).unwrap_or_default() } /// Handler for GET /api/v1/users @@ -262,17 +267,25 @@ pub async fn get_role( ) -> Result { let userid = userid.into_inner(); let tenant_id = get_tenant_id_from_request(&req); + Ok(web::Json(get_role_internal(&userid, &tenant_id)?)) +} + +/// Shared user-role implementation for HTTP handlers and in-process callers. +pub fn get_role_internal( + userid: &str, + tenant_id: &Option, +) -> Result { let tenant = tenant_id.as_deref().unwrap_or(DEFAULT_TENANT); - if !Users.contains(&userid, &tenant_id) { + if !Users.contains(userid, tenant_id) { return Err(RBACError::UserDoesNotExist); } - if let Some(p) = Users.is_protected(&userid, &tenant_id) + if let Some(p) = Users.is_protected(userid, tenant_id) && p { return Err(RBACError::ProtectedUser); }; let direct_roles: HashMap = Users - .get_role(&userid, &tenant_id) + .get_role(userid, tenant_id) .iter() .filter_map(|role_name| { if let Some(roles) = roles().get(tenant) @@ -287,7 +300,7 @@ pub async fn get_role( let mut group_roles: HashMap> = HashMap::new(); // user might be part of some user groups, fetch the roles from there as well - for user_group in Users.get_user_groups(&userid, &tenant_id) { + for user_group in Users.get_user_groups(userid, tenant_id) { if let Some(groups) = read_user_groups().get(tenant) && let Some(group) = groups.get(&user_group) { @@ -311,7 +324,7 @@ pub async fn get_role( direct_roles, group_roles, }; - Ok(web::Json(res)) + Ok(serde_json::to_value(res)?) } // Handler for DELETE /api/v1/user/delete/{userid} diff --git a/src/handlers/http/role.rs b/src/handlers/http/role.rs index dfeed50bc..969020c24 100644 --- a/src/handlers/http/role.rs +++ b/src/handlers/http/role.rs @@ -99,19 +99,31 @@ pub async fn put( pub async fn get(req: HttpRequest, name: web::Path) -> Result { let name = name.into_inner(); let tenant_id = get_tenant_id_from_request(&req); - let metadata = get_metadata(&tenant_id).await?; - let role = metadata.roles.get(&name).cloned().unwrap_or_default(); + Ok(web::Json(get_internal(&name, &tenant_id).await?)) +} + +/// Shared role lookup implementation for HTTP handlers and in-process callers. +pub async fn get_internal(name: &str, tenant_id: &Option) -> Result { + let metadata = get_metadata(tenant_id).await?; + let role = metadata.roles.get(name).cloned().unwrap_or_default(); if role.role_type().eq(&RoleType::Internal) { return Err(RoleError::ProtectedRole); } - Ok(web::Json(role)) + Ok(role) } // Handler for GET /api/v1/roles // Fetch all roles in the system pub async fn list(req: HttpRequest) -> Result { let tenant_id = get_tenant_id_from_request(&req); - let metadata = get_metadata(&tenant_id).await?; + Ok(web::Json(list_internal(&tenant_id).await?)) +} + +/// Shared role-list implementation for HTTP handlers and in-process callers. +pub async fn list_internal( + tenant_id: &Option, +) -> Result, RoleError> { + let metadata = get_metadata(tenant_id).await?; let mut roles = HashMap::new(); for (k, r) in metadata.roles.into_iter() { if !r.role_type().eq(&RoleType::Internal) { @@ -119,7 +131,7 @@ pub async fn list(req: HttpRequest) -> Result { } } - Ok(web::Json(roles)) + Ok(roles) } // Handler for DELETE /api/v1/role/{name} @@ -188,6 +200,11 @@ pub async fn put_default( // Delete existing role pub async fn get_default(req: HttpRequest) -> Result { let tenant_id = get_tenant_id_from_request(&req); + Ok(web::Json(get_default_internal(&tenant_id))) +} + +/// Shared default-role implementation for HTTP handlers and in-process callers. +pub fn get_default_internal(tenant_id: &Option) -> serde_json::Value { let tenant_id = tenant_id.as_deref().unwrap_or(DEFAULT_TENANT); let res = if let Some(role) = DEFAULT_ROLE .read() @@ -208,7 +225,7 @@ pub async fn get_default(req: HttpRequest) -> Result // None => serde_json::Value::Null, // }; - Ok(web::Json(res)) + res } async fn get_metadata( diff --git a/src/handlers/http/targets.rs b/src/handlers/http/targets.rs index 96e440d67..0d45c0156 100644 --- a/src/handlers/http/targets.rs +++ b/src/handlers/http/targets.rs @@ -36,25 +36,40 @@ use crate::{ // POST /targets pub async fn post( req: HttpRequest, - Json(mut target): Json, + Json(target): Json, ) -> Result { let tenant_id = get_tenant_id_from_request(&req); - target.tenant = tenant_id; + Ok(web::Json(post_internal(target, &tenant_id).await?)) +} + +/// Shared target-creation implementation for HTTP handlers and in-process callers. +pub async fn post_internal( + mut target: Target, + tenant_id: &Option, +) -> Result { + target.tenant = tenant_id.clone(); target.validate_outbound_policy().await?; // should check for duplicacy and liveness (??) // add to the map TARGETS.update(target.clone()).await?; - Ok(web::Json(target.mask())) + Ok(serde_json::to_value(target.mask())?) } // GET /targets pub async fn list(req: HttpRequest) -> Result { let tenant_id = get_tenant_id_from_request(&req); + Ok(web::Json(list_internal(&tenant_id).await?)) +} + +/// Shared target-list implementation for HTTP handlers and in-process callers. +pub async fn list_internal( + tenant_id: &Option, +) -> Result, AlertError> { let handles = FuturesUnordered::new(); // add to the map let mut list = vec![]; - for target in TARGETS.list(&tenant_id).await? { + for target in TARGETS.list(tenant_id).await? { handles.push(tokio::spawn(async move { if let Err(err) = target.validate_outbound_policy().await { error!(error = %err, "Alert target rejected during outbound policy validation"); @@ -82,14 +97,22 @@ pub async fn list(req: HttpRequest) -> Result { ); } - Ok(web::Json(list)) + Ok(list) } // GET /targets/{target_id} pub async fn get(req: HttpRequest, target_id: Path) -> Result { let target_id = target_id.into_inner(); let tenant_id = get_tenant_id_from_request(&req); - let target = TARGETS.get_target_by_id(&target_id, &tenant_id).await?; + Ok(web::Json(get_internal(target_id, &tenant_id).await?)) +} + +/// Shared target lookup implementation for HTTP handlers and in-process callers. +pub async fn get_internal( + target_id: Ulid, + tenant_id: &Option, +) -> Result { + let target = TARGETS.get_target_by_id(&target_id, tenant_id).await?; let res = if let Err(e) = target.validate_outbound_policy().await { let message = match &e { AlertError::OutboundPolicy(policy_error) => policy_error.sanitized_message(), @@ -106,7 +129,7 @@ pub async fn get(req: HttpRequest, target_id: Path) -> Result, ) -> Result { + let tenant_id = get_tenant_id_from_request(&req); + let target = query_target(&req, &tenant_id)?; + let response = list_traces_with_target(body, &tenant_id, target).await?; + Ok(HttpResponse::Ok().json(response)) +} + +pub async fn list_traces_internal( + body: TraceListRequest, + tenant_id: &Option, + session_key: &SessionKey, + query_auth: Option, +) -> Result { + let target = internal_query_target(session_key, query_auth)?; + list_traces_with_target(body, tenant_id, target).await +} + +async fn list_traces_with_target( + body: TraceListRequest, + tenant_id: &Option, + target: TraceQueryTarget, +) -> Result { let limit = body.limit.unwrap_or(DEFAULT_TRACE_LIMIT); if limit == 0 || limit > MAX_TRACE_LIMIT { return Err(TraceError::BadRequest(format!( @@ -201,13 +222,12 @@ pub async fn list_traces( )); } - let tenant_id = get_tenant_id_from_request(&req); - create_streams_for_distributed(vec![body.dataset.clone()], &tenant_id) + create_streams_for_distributed(vec![body.dataset.clone()], tenant_id) .await .map_err(|error| TraceError::BadRequest(error.to_string()))?; let time_range = parse_time_range(&body.start_time, &body.end_time)?; let dataset_info = - validate_trace_dataset(&body.dataset, &tenant_id, TRACE_LIST_REQUIRED_FIELDS)?; + validate_trace_dataset(&body.dataset, tenant_id, TRACE_LIST_REQUIRED_FIELDS)?; let context = TraceSqlContext::new( &body.dataset, &dataset_info.time_column, @@ -217,7 +237,6 @@ pub async fn list_traces( let conditions = build_conditions_filter(body.conditions.as_ref())?; let option = body.options.unwrap_or_default(); let sort_by = body.sort_by.unwrap_or_default(); - let target = query_target(&req, &tenant_id)?; let start_time = time_range.start.to_rfc3339(); let end_time = time_range.end.to_rfc3339(); @@ -228,7 +247,7 @@ pub async fn list_traces( build_trace_list_sql(&context, &conditions, option, sort_by, offset, limit), &start_time, &end_time, - &tenant_id, + tenant_id, ), execute_trace_query( "traces/list/count", @@ -236,7 +255,7 @@ pub async fn list_traces( build_trace_count_sql(&context, &conditions, option), &start_time, &end_time, - &tenant_id, + tenant_id, ), )?; let count = count_records @@ -245,31 +264,51 @@ pub async fn list_traces( .and_then(json_u64) .unwrap_or(0); - Ok(HttpResponse::Ok().json(TraceListResponse { + Ok(serde_json::to_value(TraceListResponse { count, offset, limit, records, - })) + }) + .expect("trace list response is serializable")) } pub async fn get_trace_detail( req: HttpRequest, Json(body): Json, ) -> Result { + let tenant_id = get_tenant_id_from_request(&req); + let target = query_target(&req, &tenant_id)?; + let response = get_trace_detail_with_target(body, &tenant_id, target).await?; + Ok(HttpResponse::Ok().json(response)) +} + +pub async fn get_trace_detail_internal( + body: TraceDetailRequest, + tenant_id: &Option, + session_key: &SessionKey, + query_auth: Option, +) -> Result { + let target = internal_query_target(session_key, query_auth)?; + get_trace_detail_with_target(body, tenant_id, target).await +} + +async fn get_trace_detail_with_target( + body: TraceDetailRequest, + tenant_id: &Option, + target: TraceQueryTarget, +) -> Result { let trace_id = body.trace_id.trim(); if trace_id.is_empty() { return Err(TraceError::BadRequest("traceId is required".to_string())); } - let tenant_id = get_tenant_id_from_request(&req); - create_streams_for_distributed(vec![body.dataset.clone()], &tenant_id) + create_streams_for_distributed(vec![body.dataset.clone()], tenant_id) .await .map_err(|error| TraceError::BadRequest(error.to_string()))?; let discovery_range = parse_time_range(&body.start_time, &body.end_time)?; let dataset_info = - validate_trace_dataset(&body.dataset, &tenant_id, TRACE_DETAIL_REQUIRED_FIELDS)?; - let target = query_target(&req, &tenant_id)?; + validate_trace_dataset(&body.dataset, tenant_id, TRACE_DETAIL_REQUIRED_FIELDS)?; let bounds = execute_trace_query( "traces/detail/bounds", @@ -277,7 +316,7 @@ pub async fn get_trace_detail( build_trace_bounds_sql(&body.dataset, trace_id, &dataset_info.time_column), &discovery_range.start.to_rfc3339(), &discovery_range.end.to_rfc3339(), - &tenant_id, + tenant_id, ) .await?; let bounds = bounds.first().ok_or_else(|| { @@ -313,15 +352,16 @@ pub async fn get_trace_detail( ), &start_time.to_rfc3339(), &(end_time + Duration::minutes(1)).to_rfc3339(), - &tenant_id, + tenant_id, ) .await?; - Ok(HttpResponse::Ok().json(TraceDetailResponse { + Ok(serde_json::to_value(TraceDetailResponse { start_time: start_time.to_rfc3339(), end_time: end_time.to_rfc3339(), records, - })) + }) + .expect("trace detail response is serializable")) } fn parse_time_range(start_time: &str, end_time: &str) -> Result { @@ -663,6 +703,21 @@ fn query_target( } } +fn internal_query_target( + session_key: &SessionKey, + query_auth: Option, +) -> Result { + match PARSEABLE.options.mode { + Mode::All | Mode::Query => Ok(TraceQueryTarget::Local(session_key.clone())), + Mode::Prism => query_auth.map(TraceQueryTarget::Remote).ok_or_else(|| { + TraceError::Internal("Missing query-node authentication headers".to_string()) + }), + mode => Err(TraceError::BadRequest(format!( + "Trace queries are not available in {mode:?} mode" + ))), + } +} + async fn execute_trace_query( query_name: &'static str, target: TraceQueryTarget, diff --git a/src/handlers/http/users/dashboards.rs b/src/handlers/http/users/dashboards.rs index 4eaab2169..8ad3df1db 100644 --- a/src/handlers/http/users/dashboards.rs +++ b/src/handlers/http/users/dashboards.rs @@ -32,6 +32,7 @@ use actix_web::{ web::{self, Json, Path}, }; use serde_json::Error as SerdeError; +use serde_json::Value; pub async fn list_dashboards(req: HttpRequest) -> Result { let tenant_id = get_tenant_id_from_request(&req); @@ -56,39 +57,52 @@ pub async fn list_dashboards(req: HttpRequest) -> Result>(); + let dashboard_summaries = list_dashboards_internal(0, Some(tags), &tenant_id).await; return Ok((web::Json(dashboard_summaries), StatusCode::OK)); } } - let dashboards = DASHBOARDS - .list_dashboards(dashboard_limit, &tenant_id) - .await; - let dashboard_summaries = dashboards - .iter() - .map(|dashboard| dashboard.to_summary()) - .collect::>(); + let dashboard_summaries = list_dashboards_internal(dashboard_limit, None, &tenant_id).await; Ok((web::Json(dashboard_summaries), StatusCode::OK)) } +pub async fn list_dashboards_internal( + limit: usize, + tags: Option>, + tenant_id: &Option, +) -> Vec { + let dashboards = if let Some(tags) = tags { + DASHBOARDS.list_dashboards_by_tags(tags, tenant_id).await + } else { + DASHBOARDS.list_dashboards(limit, tenant_id).await + }; + dashboards + .iter() + .map(|dashboard| Value::Object(dashboard.to_summary())) + .collect() +} + pub async fn get_dashboard( req: HttpRequest, dashboard_id: Path, ) -> Result { let dashboard_id = validate_dashboard_id(dashboard_id.into_inner())?; let tenant_id = get_tenant_id_from_request(&req); - let dashboard = DASHBOARDS - .get_dashboard(dashboard_id, &tenant_id) - .await - .ok_or_else(|| DashboardError::Metadata("Dashboard does not exist"))?; + let dashboard = get_dashboard_internal(dashboard_id, &tenant_id).await?; Ok((web::Json(dashboard), StatusCode::OK)) } +pub async fn get_dashboard_internal( + dashboard_id: ulid::Ulid, + tenant_id: &Option, +) -> Result { + DASHBOARDS + .get_dashboard(dashboard_id, tenant_id) + .await + .ok_or(DashboardError::Metadata("Dashboard does not exist")) +} + pub async fn create_dashboard( req: HttpRequest, Json(mut dashboard): Json, diff --git a/src/utils/mod.rs b/src/utils/mod.rs index 13fabdd85..bbf2de6e2 100644 --- a/src/utils/mod.rs +++ b/src/utils/mod.rs @@ -234,19 +234,23 @@ pub fn is_admin(req: &HttpRequest) -> Result { let session_key = extract_session_key_from_req(req).map_err(|e| anyhow::Error::msg(e.to_string()))?; - let permissions = Users.get_permissions(&session_key); + Ok(is_admin_for_session(&session_key)) +} + +pub fn is_admin_for_session(session_key: &SessionKey) -> bool { + let permissions = Users.get_permissions(session_key); // Check if user has admin permissions (Action::All on All resources) for permission in permissions.iter() { match permission { Permission::Resource(Action::All, Some(ParseableResourceType::All)) => { - return Ok(true); + return true; } _ => continue, } } - Ok(false) + false } pub fn create_intracluster_auth_headermap( From 63bcea9b6cff795f3a15d2c3685c7b3b84697618 Mon Sep 17 00:00:00 2001 From: Nikhil Sinha Date: Fri, 2 Oct 2026 11:26:43 +0530 Subject: [PATCH 2/5] fix clippy and deepsource issue --- src/handlers/http/targets.rs | 2 +- src/lib.rs | 5 +++++ 2 files changed, 6 insertions(+), 1 deletion(-) diff --git a/src/handlers/http/targets.rs b/src/handlers/http/targets.rs index 0d45c0156..ae9d312ea 100644 --- a/src/handlers/http/targets.rs +++ b/src/handlers/http/targets.rs @@ -47,7 +47,7 @@ pub async fn post_internal( mut target: Target, tenant_id: &Option, ) -> Result { - target.tenant = tenant_id.clone(); + target.tenant.clone_from(tenant_id); target.validate_outbound_policy().await?; // should check for duplicacy and liveness (??) // add to the map diff --git a/src/lib.rs b/src/lib.rs index d2e0856c1..d7439b4db 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,3 +1,8 @@ +#![allow(clippy::double_must_use)] +// `async_trait` currently expands async trait methods with a redundant +// `must_use` marker under Rust 1.99. Keep the compatibility allowance scoped +// to this lint until the upstream macro/compiler combination stops emitting it. + /* * Parseable Server (C) 2022 - 2025 Parseable, Inc. * From 3e2c0220619fa7ef101dc460f50d16a2645201bf Mon Sep 17 00:00:00 2001 From: Nikhil Sinha Date: Fri, 2 Oct 2026 17:26:11 +0530 Subject: [PATCH 3/5] oss endpoint for tools --- src/handlers/http/modal/query_server.rs | 1 + src/handlers/http/modal/server.rs | 11 + src/lib.rs | 2 + src/tool_catalog.rs | 353 ++++++++++++++++++++++++ 4 files changed, 367 insertions(+) create mode 100644 src/tool_catalog.rs diff --git a/src/handlers/http/modal/query_server.rs b/src/handlers/http/modal/query_server.rs index 278f7c4b6..dd2208f79 100644 --- a/src/handlers/http/modal/query_server.rs +++ b/src/handlers/http/modal/query_server.rs @@ -80,6 +80,7 @@ impl ParseableServer for QueryServer { .service( web::scope(&prism_base_path()) .service(Server::get_prism_home()) + .service(Server::get_llm_webscope()) .service(Server::get_prism_logstream()) .service(Server::get_prism_datasets()) .service(Server::get_apikeys_webscope()) diff --git a/src/handlers/http/modal/server.rs b/src/handlers/http/modal/server.rs index 686823501..e711d2ca4 100644 --- a/src/handlers/http/modal/server.rs +++ b/src/handlers/http/modal/server.rs @@ -111,6 +111,7 @@ impl ParseableServer for Server { .service( web::scope(&prism_base_path()) .service(Server::get_prism_home()) + .service(Server::get_llm_webscope()) .service(Server::get_prism_logstream()) .service(Server::get_prism_datasets()) .service(Server::get_apikeys_webscope()) @@ -199,6 +200,16 @@ impl ParseableServer for Server { } impl Server { + pub fn get_llm_webscope() -> Scope { + web::scope("/llm").service( + web::resource("/tools").route( + web::get() + .to(crate::tool_catalog::list_tool_registry) + .authorize(Action::Query), + ), + ) + } + pub fn get_prism_home() -> Scope { web::scope("/home") .service(web::resource("").route(web::get().to(http::prism_home::home_api))) diff --git a/src/lib.rs b/src/lib.rs index d7439b4db..a8006d364 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -57,6 +57,7 @@ pub mod storage; pub mod sync; pub mod telemetry; pub mod tenants; +pub mod tool_catalog; pub mod users; pub mod utils; pub mod validator; @@ -66,6 +67,7 @@ use std::time::Duration; pub use arrow_array; pub use arrow_flight; pub use arrow_ipc; +pub use async_trait; pub use catalog as parseable_catalog; pub use datafusion; pub use datafusion_proto; diff --git a/src/tool_catalog.rs b/src/tool_catalog.rs new file mode 100644 index 000000000..0d6b79141 --- /dev/null +++ b/src/tool_catalog.rs @@ -0,0 +1,353 @@ +use actix_web::HttpResponse; +use serde_json::{Value, json}; + +#[derive(Clone, Debug)] +pub struct ToolSpec { + pub name: &'static str, + pub title: &'static str, + pub description: &'static str, + pub input_schema: Value, +} + +impl ToolSpec { + fn new( + name: &'static str, + title: &'static str, + description: &'static str, + input_schema: Value, + ) -> Self { + Self { + name, + title, + description, + input_schema, + } + } +} + +fn object_schema(properties: Value, required: &[&str]) -> Value { + json!({ + "type": "object", + "properties": properties, + "required": required, + "additionalProperties": false, + }) +} + +fn empty_schema() -> Value { + object_schema(json!({}), &[]) +} + +fn string_property(description: &str) -> Value { + json!({ "type": "string", "minLength": 1, "description": description }) +} + +fn dataset_schema() -> Value { + object_schema( + json!({ "dataset": string_property("Dataset (log stream) name.") }), + &["dataset"], + ) +} + +fn alert_schema() -> Value { + object_schema( + json!({ "alert_id": string_property("Alert ID. Use list_alerts to discover.") }), + &["alert_id"], + ) +} + +fn target_schema() -> Value { + object_schema( + json!({ "target_id": string_property("Alert target ID. Use list_alert_targets to discover.") }), + &["target_id"], + ) +} + +/// Tools implemented by both Parseable OSS and Enterprise. +/// +/// Enterprise may extend this catalog, but an OSS server only advertises this +/// list so MCP clients never discover enterprise-only capabilities. +pub fn oss_tool_specs() -> Vec { + vec![ + ToolSpec::new( + "list_datasets", + "List datasets", + "List all log datasets available on this Parseable server.", + empty_schema(), + ), + ToolSpec::new( + "get_dataset_schema", + "Get dataset schema", + "Get columns and data types for a dataset.", + dataset_schema(), + ), + ToolSpec::new( + "get_dataset_info", + "Get dataset info", + "Get metadata for a dataset, including event times and storage configuration.", + dataset_schema(), + ), + ToolSpec::new( + "get_dataset_stats", + "Get dataset stats", + "Get event count, ingested bytes, and storage bytes for a dataset.", + dataset_schema(), + ), + ToolSpec::new( + "sample_events", + "Sample recent events", + "Return recent events from a dataset.", + object_schema( + json!({ + "dataset": string_property("Dataset name."), + "limit": { "type": "integer", "minimum": 1, "maximum": 100, "default": 10 }, + "minutes": { "type": "integer", "minimum": 1, "maximum": 1440, "default": 60 }, + }), + &["dataset"], + ), + ), + ToolSpec::new( + "get_log_context", + "Get surrounding log context", + "Fetch records around an exact anchor log row.", + object_schema( + json!({ + "dataset": string_property("Dataset containing the anchor row."), + "pTimestamp": string_property("Exact p_timestamp from the anchor row."), + "log": { "type": "string" }, + "body": { "type": "string" }, + "message": { "type": "string" }, + "contextWindow": { "type": "string" }, + "contextStartTime": { "type": "string" }, + "contextEndTime": { "type": "string" }, + "conditions": { "type": "object", "additionalProperties": true }, + "pageSize": { "type": "integer", "minimum": 1, "maximum": 500, "default": 100 }, + }), + &["dataset", "pTimestamp"], + ), + ), + ToolSpec::new( + "query_sql", + "Run SQL query", + "Execute a read-only SQL query over a time window.", + object_schema( + json!({ + "query": string_property("Read-only DataFusion SQL SELECT."), + "startTime": string_property("Start of time window."), + "endTime": string_property("End of time window."), + }), + &["query", "startTime", "endTime"], + ), + ), + ToolSpec::new( + "query_promql", + "Run PromQL query", + "Execute a PromQL instant or range query.", + object_schema( + json!({ + "query": string_property("PromQL expression."), + "stream": string_property("Metrics dataset name."), + "start": { "type": "string" }, + "end": { "type": "string" }, + "step": { "type": "string" }, + "time": { "type": "string" }, + "timeout": { "type": "string" }, + "limit": { "type": "string" }, + "timestamp_format": { "type": "string", "enum": ["rfc3339", "unix"] }, + }), + &["query", "stream"], + ), + ), + ToolSpec::new( + "list_alerts", + "List alerts", + "List alerts configured on this Parseable server.", + empty_schema(), + ), + ToolSpec::new( + "get_alert", + "Get alert", + "Get the full configuration for one alert.", + alert_schema(), + ), + ToolSpec::new( + "list_alert_tags", + "List alert tags", + "List tags used across alerts.", + empty_schema(), + ), + ToolSpec::new( + "enable_alert", + "Enable alert", + "Enable a disabled alert.", + alert_schema(), + ), + ToolSpec::new( + "disable_alert", + "Disable alert", + "Disable an active alert.", + alert_schema(), + ), + ToolSpec::new( + "evaluate_alert", + "Evaluate alert now", + "Evaluate an alert immediately. This may trigger notifications.", + alert_schema(), + ), + ToolSpec::new( + "create_alert", + "Create alert", + "Create an alert from a complete Parseable alert specification.", + object_schema( + json!({ "spec": { "type": "object", "additionalProperties": true } }), + &["spec"], + ), + ), + ToolSpec::new( + "list_alert_targets", + "List alert targets", + "List configured alert notification targets.", + empty_schema(), + ), + ToolSpec::new( + "get_alert_target", + "Get alert target", + "Get the configuration for one alert target.", + target_schema(), + ), + ToolSpec::new( + "create_alert_target", + "Create alert target", + "Create an alert notification target.", + object_schema( + json!({ "spec": { "type": "object", "additionalProperties": true } }), + &["spec"], + ), + ), + ToolSpec::new( + "ping", + "Ping Parseable server", + "Check connectivity and server health.", + empty_schema(), + ), + ToolSpec::new( + "explain_query", + "Explain SQL query plan", + "Explain a read-only SQL query without executing it.", + object_schema( + json!({ + "query": string_property("SQL SELECT to explain."), + "startTime": string_property("Start of time window."), + "endTime": string_property("End of time window."), + }), + &["query", "startTime", "endTime"], + ), + ), + ToolSpec::new( + "list_users", + "List users", + "List registered users.", + empty_schema(), + ), + ToolSpec::new( + "get_user_roles", + "Get user roles", + "Get roles assigned to a user.", + object_schema( + json!({ "userid": string_property("User ID.") }), + &["userid"], + ), + ), + ToolSpec::new( + "list_roles", + "List roles", + "List roles defined on this Parseable server.", + empty_schema(), + ), + ToolSpec::new( + "get_role", + "Get role", + "Get the privileges assigned to a role.", + object_schema(json!({ "name": string_property("Role name.") }), &["name"]), + ), + ToolSpec::new( + "get_default_role", + "Get default role", + "Get the role assigned to new users by default.", + empty_schema(), + ), + ToolSpec::new( + "get_cluster_status", + "Get cluster status", + "List nodes and their status in a distributed deployment.", + empty_schema(), + ), + ToolSpec::new( + "get_cluster_metrics", + "Get cluster metrics", + "Get aggregated metrics for a distributed deployment.", + empty_schema(), + ), + ToolSpec::new( + "get_retention", + "Get dataset retention policy", + "Get the retention policy for a dataset.", + dataset_schema(), + ), + ] +} + +pub fn is_oss_tool(name: &str) -> bool { + oss_tool_specs().iter().any(|tool| tool.name == name) +} + +/// `GET /api/prism/v1/llm/tools` +pub async fn list_tool_registry() -> HttpResponse { + let tools = oss_tool_specs() + .into_iter() + .map(|tool| { + json!({ + "name": tool.name, + "title": tool.title, + "description": tool.description, + "inputSchema": tool.input_schema, + }) + }) + .collect::>(); + + HttpResponse::Ok().json(json!({ "tools": tools })) +} + +#[cfg(test)] +mod tests { + use std::collections::HashSet; + + use actix_web::{App, test, web}; + use serde_json::Value; + + use super::{list_tool_registry, oss_tool_specs}; + + #[actix_web::test] + async fn oss_catalog_has_unique_names() { + let tools = oss_tool_specs(); + let names = tools.iter().map(|tool| tool.name).collect::>(); + + assert_eq!(tools.len(), 28); + assert_eq!(tools.len(), names.len()); + } + + #[actix_web::test] + async fn registry_api_returns_only_oss_tools() { + let app = + test::init_service(App::new().route("/tools", web::get().to(list_tool_registry))).await; + let response: Value = test::call_and_read_body_json( + &app, + test::TestRequest::get().uri("/tools").to_request(), + ) + .await; + + assert_eq!(response["tools"].as_array().unwrap().len(), 28); + assert_eq!(response["tools"][0]["name"], "list_datasets"); + assert_eq!(response["tools"][0]["inputSchema"]["type"], "object"); + } +} From e8f15f2cab93ff8163d333a636d19247c0934917 Mon Sep 17 00:00:00 2001 From: Nikhil Sinha Date: Sat, 3 Oct 2026 11:14:59 +0530 Subject: [PATCH 4/5] limit and timeout as integer --- src/tool_catalog.rs | 18 ++++++++++++++++-- 1 file changed, 16 insertions(+), 2 deletions(-) diff --git a/src/tool_catalog.rs b/src/tool_catalog.rs index 0d6b79141..e341d1dd4 100644 --- a/src/tool_catalog.rs +++ b/src/tool_catalog.rs @@ -151,8 +151,8 @@ pub fn oss_tool_specs() -> Vec { "end": { "type": "string" }, "step": { "type": "string" }, "time": { "type": "string" }, - "timeout": { "type": "string" }, - "limit": { "type": "string" }, + "timeout": { "type": "number", "minimum": 0 }, + "limit": { "type": "integer", "minimum": 1, "maximum": 500, "default": 500 }, "timestamp_format": { "type": "string", "enum": ["rfc3339", "unix"] }, }), &["query", "stream"], @@ -350,4 +350,18 @@ mod tests { assert_eq!(response["tools"][0]["name"], "list_datasets"); assert_eq!(response["tools"][0]["inputSchema"]["type"], "object"); } + + #[actix_web::test] + async fn promql_execution_limits_use_numeric_schema_types() { + let tools = oss_tool_specs(); + let promql = tools + .iter() + .find(|tool| tool.name == "query_promql") + .unwrap(); + let properties = &promql.input_schema["properties"]; + + assert_eq!(properties["timeout"]["type"], "number"); + assert_eq!(properties["limit"]["type"], "integer"); + assert_eq!(properties["limit"]["maximum"], 500); + } } From 1519ca4ec6310d45c4f1bb5e8d99751349897523 Mon Sep 17 00:00:00 2001 From: Nikhil Sinha Date: Sat, 3 Oct 2026 18:11:25 +0530 Subject: [PATCH 5/5] add tools to tool catalog --- src/handlers/http/logstream.rs | 31 +++++++++++---------- src/handlers/http/rbac.rs | 14 +++++++--- src/handlers/http/users/dashboards.rs | 8 ++++-- src/handlers/http/users/filters.rs | 40 +++++++++++++++++++-------- src/tool_catalog.rs | 40 +++++++++++++++++++++++++-- 5 files changed, 98 insertions(+), 35 deletions(-) diff --git a/src/handlers/http/logstream.rs b/src/handlers/http/logstream.rs index 6e427e8a3..46f7df355 100644 --- a/src/handlers/http/logstream.rs +++ b/src/handlers/http/logstream.rs @@ -527,24 +527,27 @@ pub async fn get_stream_hot_tier( tracing::Span::current() .record("stream", tracing::field::display(&stream_name)) .record("tenant", tracing::field::debug(&tenant_id)); - // For query mode, if the stream not found in memory map, - //check if it exists in the storage - //create stream and schema from storage - if !PARSEABLE - .check_or_load_stream(&stream_name, &tenant_id) - .await - { - return Err(StreamNotFound(stream_name.clone()).into()); + let meta = get_stream_hot_tier_internal(&stream_name, &tenant_id).await?; + + Ok((web::Json(meta), StatusCode::OK)) +} + +pub async fn get_stream_hot_tier_internal( + stream_name: &str, + tenant_id: &Option, +) -> Result { + // For query mode, if the stream is not in memory, load it from storage. + if !PARSEABLE.check_or_load_stream(stream_name, tenant_id).await { + return Err(StreamNotFound(stream_name.to_owned()).into()); } let Some(hot_tier_manager) = GLOBAL_HOTTIER.get() else { - return Err(StreamError::HotTierNotEnabled(stream_name)); + return Err(StreamError::HotTierNotEnabled(stream_name.to_owned())); }; - let meta = hot_tier_manager - .get_hot_tier(&stream_name, &tenant_id) - .await?; - - Ok((web::Json(meta), StatusCode::OK)) + hot_tier_manager + .get_hot_tier(stream_name, tenant_id) + .await + .map_err(Into::into) } #[tracing::instrument( diff --git a/src/handlers/http/rbac.rs b/src/handlers/http/rbac.rs index c1dcab41d..6d5bf1896 100644 --- a/src/handlers/http/rbac.rs +++ b/src/handlers/http/rbac.rs @@ -110,15 +110,21 @@ pub async fn get_prism_user( ) -> Result { let userid = userid.into_inner(); let tenant_id = get_tenant_id_from_request(&req); + Ok(web::Json(get_user_internal(&userid, &tenant_id)?)) +} + +/// Shared user lookup implementation for HTTP handlers and in-process callers. +pub fn get_user_internal( + userid: &str, + tenant_id: &Option, +) -> Result { // First check if the user exists let users = rbac::map::users(); if let Some(users) = users.get(tenant_id.as_deref().unwrap_or(DEFAULT_TENANT)) - && let Some(user) = users.get(&userid) + && let Some(user) = users.get(userid) && !user.protected { - // Create UsersPrism for the found user only - let prism_user = to_prism_user(user); - Ok(web::Json(prism_user)) + Ok(to_prism_user(user)) } else { Err(RBACError::UserDoesNotExist) } diff --git a/src/handlers/http/users/dashboards.rs b/src/handlers/http/users/dashboards.rs index 8ad3df1db..d7addb905 100644 --- a/src/handlers/http/users/dashboards.rs +++ b/src/handlers/http/users/dashboards.rs @@ -259,12 +259,14 @@ pub async fn add_tile( } pub async fn list_tags(req: HttpRequest) -> Result { - let tags = DASHBOARDS - .list_tags(&get_tenant_id_from_request(&req)) - .await; + let tags = list_tags_internal(&get_tenant_id_from_request(&req)).await; Ok((web::Json(tags), StatusCode::OK)) } +pub async fn list_tags_internal(tenant_id: &Option) -> Vec { + DASHBOARDS.list_tags(tenant_id).await +} + #[derive(Debug, thiserror::Error)] pub enum DashboardError { #[error("Failed to connect to storage: {0}")] diff --git a/src/handlers/http/users/filters.rs b/src/handlers/http/users/filters.rs index 22976fd9b..c7a4d5681 100644 --- a/src/handlers/http/users/filters.rs +++ b/src/handlers/http/users/filters.rs @@ -20,10 +20,12 @@ use crate::{ handlers::http::rbac::RBACError, metastore::MetastoreError, parseable::PARSEABLE, + rbac::{Users, map::SessionKey}, storage::{ObjectStorageError, StreamType}, users::filters::{CURRENT_FILTER_VERSION, FILTERS, Filter}, utils::{ actix::extract_session_key_from_req, get_hash, get_user_and_tenant_from_request, is_admin, + is_admin_for_session, }, validator, }; @@ -39,27 +41,41 @@ use ulid::Ulid; pub async fn list(req: HttpRequest) -> Result { let key = extract_session_key_from_req(&req).map_err(|e| FiltersError::Custom(e.to_string()))?; - let filters = FILTERS.list_filters(&key).await; + let filters = list_internal(&key).await; Ok((web::Json(filters), StatusCode::OK)) } +pub async fn list_internal(key: &SessionKey) -> Vec { + FILTERS.list_filters(key).await +} + pub async fn get( req: HttpRequest, filter_id: Path, ) -> Result { - let (user_id, tenant_id) = get_user_and_tenant_from_request(&req)?; + let key = + extract_session_key_from_req(&req).map_err(|e| FiltersError::Custom(e.to_string()))?; + let (_, tenant_id) = get_user_and_tenant_from_request(&req)?; let filter_id = filter_id.into_inner(); - let is_admin = is_admin(&req).map_err(|e| FiltersError::Custom(e.to_string()))?; - if let Some(filter) = FILTERS - .get_filter(&filter_id, &get_hash(&user_id), is_admin, &tenant_id) - .await - { - return Ok((web::Json(filter), StatusCode::OK)); - } + let filter = get_internal(&filter_id, &key, &tenant_id).await?; + Ok((web::Json(filter), StatusCode::OK)) +} - Err(FiltersError::Metadata( - "Filter does not exist or user is not authorized", - )) +pub async fn get_internal( + filter_id: &str, + key: &SessionKey, + tenant_id: &Option, +) -> Result { + let user_id = Users + .get_userid_from_session(key) + .map(|(user_id, _)| get_hash(&user_id)) + .ok_or_else(|| FiltersError::Custom("unknown user session".to_owned()))?; + FILTERS + .get_filter(filter_id, &user_id, is_admin_for_session(key), tenant_id) + .await + .ok_or(FiltersError::Metadata( + "Filter does not exist or user is not authorized", + )) } pub async fn post( diff --git a/src/tool_catalog.rs b/src/tool_catalog.rs index e341d1dd4..1e24d0122 100644 --- a/src/tool_catalog.rs +++ b/src/tool_catalog.rs @@ -249,6 +249,15 @@ pub fn oss_tool_specs() -> Vec { "List registered users.", empty_schema(), ), + ToolSpec::new( + "get_user", + "Get user", + "Get one user's profile, direct roles, group roles, and group membership.", + object_schema( + json!({ "userid": string_property("Exact user ID.") }), + &["userid"], + ), + ), ToolSpec::new( "get_user_roles", "Get user roles", @@ -294,6 +303,33 @@ pub fn oss_tool_specs() -> Vec { "Get the retention policy for a dataset.", dataset_schema(), ), + ToolSpec::new( + "get_hot_tier_config", + "Get hot-tier configuration", + "Get hot-tier configuration and utilization for a dataset.", + dataset_schema(), + ), + ToolSpec::new( + "list_filters", + "List saved filters", + "List saved filters visible to the caller.", + empty_schema(), + ), + ToolSpec::new( + "get_filter", + "Get saved filter", + "Get one saved filter by ID when visible to the caller.", + object_schema( + json!({ "id": string_property("Saved-filter ID.") }), + &["id"], + ), + ), + ToolSpec::new( + "list_dashboard_tags", + "List dashboard tags", + "List unique dashboard tags in the caller's tenant.", + empty_schema(), + ), ] } @@ -332,7 +368,7 @@ mod tests { let tools = oss_tool_specs(); let names = tools.iter().map(|tool| tool.name).collect::>(); - assert_eq!(tools.len(), 28); + assert_eq!(tools.len(), 33); assert_eq!(tools.len(), names.len()); } @@ -346,7 +382,7 @@ mod tests { ) .await; - assert_eq!(response["tools"].as_array().unwrap().len(), 28); + assert_eq!(response["tools"].as_array().unwrap().len(), 33); assert_eq!(response["tools"][0]["name"], "list_datasets"); assert_eq!(response["tools"][0]["inputSchema"]["type"], "object"); }