fix: enforce permissions for restricted HTTP endpoints (#8672)

* chore: add perm check

Signed-off-by: shuiyisong <xixing.sys@gmail.com>

* fix: CR issues

Signed-off-by: shuiyisong <xixing.sys@gmail.com>

* chore: add comments

Signed-off-by: shuiyisong <xixing.sys@gmail.com>

---------

Signed-off-by: shuiyisong <xixing.sys@gmail.com>
This commit is contained in:
shuiyisong
2026-07-29 08:51:17 +00:00
committed by GitHub
parent 1a24a6e782
commit 2afcd4a558
7 changed files with 237 additions and 5 deletions
+31 -3
View File
@@ -105,6 +105,11 @@ pub enum PermissionReq<'a> {
PromStoreRead,
Otlp,
LogWrite,
JaegerQuery,
PipelineQuery,
PipelineManage,
DashboardQuery,
DashboardManage,
BulkInsert {
catalog: &'a str,
schema: &'a str,
@@ -119,9 +124,12 @@ impl<'a> PermissionReq<'a> {
PermissionReq::GrpcRequest(Request::Query(query_request)) => {
!matches!(query_request.query, Some(Query::InsertIntoPlan(_)))
}
PermissionReq::PromQuery | PermissionReq::LogQuery | PermissionReq::PromStoreRead => {
true
}
PermissionReq::PromQuery
| PermissionReq::LogQuery
| PermissionReq::PromStoreRead
| PermissionReq::JaegerQuery
| PermissionReq::PipelineQuery
| PermissionReq::DashboardQuery => true,
PermissionReq::SqlStatement(stmt) => stmt.is_readonly(),
PermissionReq::GrpcRequest(_)
@@ -130,6 +138,8 @@ impl<'a> PermissionReq<'a> {
| PermissionReq::PromStoreWrite
| PermissionReq::Otlp
| PermissionReq::LogWrite
| PermissionReq::PipelineManage
| PermissionReq::DashboardManage
| PermissionReq::BulkInsert { .. } => false,
}
}
@@ -489,6 +499,24 @@ mod tests {
assert!(req.is_write());
}
#[test]
fn test_management_request_access_modes() {
for req in [
PermissionReq::JaegerQuery,
PermissionReq::PipelineQuery,
PermissionReq::DashboardQuery,
] {
assert!(req.is_readonly());
}
for req in [
PermissionReq::PipelineManage,
PermissionReq::DashboardManage,
] {
assert!(req.is_write());
}
}
#[test]
fn test_table_target_permission_forwarding() {
let checker: PermissionCheckerRef = Arc::new(TargetAwarePermissionChecker);
+158 -1
View File
@@ -161,6 +161,19 @@ impl Instance {
&self.plugins
}
fn check_permission(
&self,
ctx: &QueryContextRef,
req: PermissionReq<'_>,
) -> server_error::Result<()> {
self.plugins
.get::<PermissionCheckerRef>()
.as_ref()
.check_permission(ctx.current_user(), req)
.context(AuthSnafu)?;
Ok(())
}
pub fn statement_executor(&self) -> &StatementExecutorRef {
&self.statement_executor
}
@@ -1682,7 +1695,10 @@ mod tests {
use datatypes::vectors::{StringVector, VectorRef};
use log_query::LogQuery;
use query::query_engine::options::QueryOptions;
use servers::query_handler::{LogQueryHandler, PromStoreProtocolHandler};
use servers::query_handler::{
DashboardHandler, JaegerQueryHandler, LogQueryHandler, PipelineHandler,
PromStoreProtocolHandler,
};
use session::context::{Channel, ConnInfo, QueryContext, QueryContextBuilder};
use snafu::{Location, Snafu};
use sql::dialect::GreptimeDbDialect;
@@ -1697,6 +1713,7 @@ mod tests {
use table::test_util::{EmptyTable, MemTable};
use table::{Table, TableRef};
use tokio::sync::{mpsc, oneshot};
use tower::ServiceExt;
use super::*;
use crate::frontend::FrontendOptions;
@@ -1935,6 +1952,54 @@ mod tests {
}
}
struct RejectEndpointPermissionChecker;
impl PermissionChecker for RejectEndpointPermissionChecker {
fn check_permission(
&self,
_user_info: UserInfoRef,
req: PermissionReq,
) -> auth::error::Result<PermissionResp> {
Ok(
if matches!(
req,
PermissionReq::PipelineQuery
| PermissionReq::PipelineManage
| PermissionReq::DashboardQuery
| PermissionReq::DashboardManage
) {
PermissionResp::Reject
} else {
PermissionResp::Allow
},
)
}
fn check_permission_with_table_targets(
&self,
user_info: UserInfoRef,
req: PermissionReq,
targets: PermissionTableTargets,
) -> auth::error::Result<PermissionResp> {
match req {
PermissionReq::JaegerQuery => {
let reject = match targets {
PermissionTableTargets::Unresolved => true,
PermissionTableTargets::Resolved(targets) => {
targets.iter().any(|target| target.table == "denied")
}
};
Ok(if reject {
PermissionResp::Reject
} else {
PermissionResp::Allow
})
}
req => self.check_permission(user_info, req),
}
}
}
#[derive(Default)]
struct TargetIndependentPermissionChecker {
checks: atomic::AtomicUsize,
@@ -2412,6 +2477,14 @@ mod tests {
})
}
fn assert_permission_denied<T>(result: servers::error::Result<T>) {
let err = match result {
Ok(_) => panic!("request should be rejected"),
Err(err) => err,
};
assert_eq!(StatusCode::PermissionDenied, err.status_code());
}
#[tokio::test]
async fn test_event_recorder_is_exposed() -> TestResult<()> {
let instance =
@@ -2423,6 +2496,90 @@ mod tests {
Ok(())
}
#[tokio::test]
async fn test_restricted_endpoint_handlers_check_permissions() -> TestResult<()> {
let plugins = Plugins::new();
plugins.insert::<PermissionCheckerRef>(Arc::new(RejectEndpointPermissionChecker));
let instance = test_instance_with_plugins(
test_table(1024, "denied")?,
test_table(1025, "target")?,
plugins,
)
.await?;
let mut ctx = test_query_ctx(1);
Arc::get_mut(&mut ctx).unwrap().set_extension(
servers::http::jaeger::JAEGER_QUERY_TABLE_NAME_KEY,
"denied".to_string(),
);
assert_permission_denied(JaegerQueryHandler::get_services(&instance, ctx.clone()).await);
assert_permission_denied(
JaegerQueryHandler::get_operations(&instance, ctx.clone(), "service", None).await,
);
assert_permission_denied(
JaegerQueryHandler::get_trace(&instance, ctx.clone(), "trace", None, None, None).await,
);
assert_permission_denied(
JaegerQueryHandler::find_traces(
&instance,
ctx.clone(),
servers::http::jaeger::QueryTraceParams {
service_name: "service".to_string(),
..Default::default()
},
)
.await,
);
assert_permission_denied(
PipelineHandler::get_pipeline_str(&instance, "pipeline", None, ctx.clone()).await,
);
assert_permission_denied(
PipelineHandler::insert_pipeline(
&instance,
"pipeline",
"application/yaml",
"",
ctx.clone(),
)
.await,
);
assert_permission_denied(
PipelineHandler::delete_pipeline(&instance, "pipeline", None, ctx.clone()).await,
);
let app = axum::Router::new()
.route(
"/pipelines/_dryrun",
axum::routing::post(servers::http::event::pipeline_dryrun),
)
.with_state(servers::http::event::LogState {
log_handler: Arc::new(instance.clone()),
log_validator: None,
ingest_interceptor: None,
})
.layer(axum::Extension((*ctx).clone()));
let response = app
.oneshot(
axum::http::Request::post("/pipelines/_dryrun")
.header("content-type", "application/json")
.body(axum::body::Body::from("{}"))
.unwrap(),
)
.await
.unwrap();
assert_eq!(axum::http::StatusCode::FORBIDDEN, response.status());
assert_permission_denied(
DashboardHandler::save(&instance, "dashboard", "{}", ctx.clone()).await,
);
assert_permission_denied(DashboardHandler::list(&instance, ctx.clone()).await);
assert_permission_denied(
DashboardHandler::delete(&instance, "dashboard", ctx.clone()).await,
);
Ok(())
}
#[tokio::test]
async fn test_target_independent_checker_skips_target_resolution() -> TestResult<()> {
let physical_table = "physical_metric";
+4
View File
@@ -21,6 +21,7 @@ use api::v1::{
RowInsertRequests, Rows, SemanticType,
};
use async_trait::async_trait;
use auth::PermissionReq;
use common_catalog::consts::{DEFAULT_PRIVATE_SCHEMA_NAME, default_engine};
use common_error::ext::BoxedError;
use common_query::OutputData;
@@ -392,14 +393,17 @@ impl servers::query_handler::DashboardHandler for Instance {
definition: &str,
ctx: QueryContextRef,
) -> servers::error::Result<()> {
self.check_permission(&ctx, PermissionReq::DashboardManage)?;
self.insert_dashboard(name, definition, ctx).await
}
async fn list(&self, ctx: QueryContextRef) -> servers::error::Result<Vec<DashboardDefinition>> {
self.check_permission(&ctx, PermissionReq::DashboardQuery)?;
self.list_dashboards(ctx).await
}
async fn delete(&self, name: &str, ctx: QueryContextRef) -> servers::error::Result<()> {
self.check_permission(&ctx, PermissionReq::DashboardManage)?;
self.delete_dashboard(name, ctx).await
}
}
+27 -1
View File
@@ -16,6 +16,7 @@ use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use async_trait::async_trait;
use auth::{PermissionReq, PermissionTableTarget, PermissionTableTargets};
use catalog::CatalogManagerRef;
use common_catalog::consts::{
TRACE_TABLE_NAME, trace_operations_table_name, trace_services_table_name,
@@ -38,7 +39,7 @@ use datafusion_expr::{Expr, ExprFunctionExt, SortExpr, col, lit, lit_timestamp_n
use query::QueryEngineRef;
use serde_json::Value as JsonValue;
use servers::error::{
CollectRecordbatchSnafu, DataFusionSnafu, Result as ServerResult, TableNotFoundSnafu,
AuthSnafu, CollectRecordbatchSnafu, DataFusionSnafu, Result as ServerResult, TableNotFoundSnafu,
};
use servers::http::jaeger::{JAEGER_QUERY_TABLE_NAME_KEY, QueryTraceParams, TraceUserAgent};
use servers::otlp::trace::{
@@ -58,9 +59,28 @@ use crate::instance::Instance;
const DEFAULT_LIMIT: usize = 2000;
const KEY_RN: &str = "greptime_rn";
impl Instance {
async fn check_jaeger_query_permission(&self, ctx: &QueryContextRef) -> ServerResult<()> {
let table = ctx
.extension(JAEGER_QUERY_TABLE_NAME_KEY)
.unwrap_or(TRACE_TABLE_NAME);
let targets = PermissionTableTargets::resolved(vec![PermissionTableTarget::new(
ctx.current_catalog(),
ctx.current_schema(),
table,
)]);
let targets = self.resolve_query_permission_targets(targets, ctx).await?;
self.check_table_permission(ctx, PermissionReq::JaegerQuery, targets)
.context(AuthSnafu)?;
Ok(())
}
}
#[async_trait]
impl JaegerQueryHandler for Instance {
async fn get_services(&self, ctx: QueryContextRef) -> ServerResult<Output> {
self.check_jaeger_query_permission(&ctx).await?;
// It's equivalent to `SELECT DISTINCT(service_name) FROM {db}.{trace_table}`.
Ok(query_trace_table(
ctx,
@@ -81,6 +101,8 @@ impl JaegerQueryHandler for Instance {
service_name: &str,
span_kind: Option<&str>,
) -> ServerResult<Output> {
self.check_jaeger_query_permission(&ctx).await?;
let mut filters = vec![col(SERVICE_NAME_COLUMN).eq(lit(service_name))];
if let Some(span_kind) = span_kind {
@@ -129,6 +151,8 @@ impl JaegerQueryHandler for Instance {
end_time: Option<i64>,
limit: Option<usize>,
) -> ServerResult<Output> {
self.check_jaeger_query_permission(&ctx).await?;
// It's equivalent to the following SQL query:
//
// ```
@@ -173,6 +197,8 @@ impl JaegerQueryHandler for Instance {
ctx: QueryContextRef,
query_params: QueryTraceParams,
) -> ServerResult<Output> {
self.check_jaeger_query_permission(&ctx).await?;
let mut filters = vec![];
// `service_name` is already validated in `from_jaeger_query_params()`, so no additional check needed here.
+7
View File
@@ -84,6 +84,10 @@ impl PipelineHandler for Instance {
Ok(outputs)
}
fn check_pipeline_query_permission(&self, query_ctx: &QueryContextRef) -> ServerResult<()> {
self.check_permission(query_ctx, PermissionReq::PipelineQuery)
}
async fn get_pipeline(
&self,
name: &str,
@@ -103,6 +107,7 @@ impl PipelineHandler for Instance {
pipeline: &str,
query_ctx: QueryContextRef,
) -> ServerResult<PipelineInfo> {
self.check_permission(&query_ctx, PermissionReq::PipelineManage)?;
self.pipeline_operator
.insert_pipeline(name, content_type, pipeline, query_ctx)
.await
@@ -115,6 +120,7 @@ impl PipelineHandler for Instance {
version: PipelineVersion,
ctx: QueryContextRef,
) -> ServerResult<Option<()>> {
self.check_permission(&ctx, PermissionReq::PipelineManage)?;
self.pipeline_operator
.delete_pipeline(name, version, ctx)
.await
@@ -143,6 +149,7 @@ impl PipelineHandler for Instance {
version: PipelineVersion,
query_ctx: QueryContextRef,
) -> ServerResult<(String, TimestampNanosecond)> {
self.check_permission(&query_ctx, PermissionReq::PipelineQuery)?;
self.pipeline_operator
.get_pipeline_str(name, version, query_ctx)
.await
+2
View File
@@ -240,6 +240,7 @@ pub async fn query_pipeline_ddl(
query_ctx.set_channel(Channel::Log);
let query_ctx = Arc::new(query_ctx);
handler.check_pipeline_query_permission(&query_ctx)?;
let pipeline = handler
.get_pipeline(&pipeline_name, version, query_ctx.clone())
.await?;
@@ -614,6 +615,7 @@ pub async fn pipeline_dryrun(
query_ctx.set_channel(Channel::Log);
let query_ctx = Arc::new(query_ctx);
handler.check_pipeline_query_permission(&query_ctx)?;
match check_pipeline_dryrun_params_valid(&payload) {
Some(params) => {
+8
View File
@@ -177,6 +177,14 @@ pub trait PipelineHandler {
inputs: Vec<(QueryContextRef, RowInsertRequests)>,
) -> Result<Vec<Result<Output>>>;
fn check_pipeline_query_permission(&self, query_ctx: &QueryContextRef) -> Result<()>;
/// Loads a compiled pipeline for execution.
///
/// This intentionally does not check pipeline-query permission: users with
/// write-only permission can ingest through an existing pipeline. Inspection
/// and preview callers must check query permission first; ingestion enforces
/// write and table-target permissions separately.
async fn get_pipeline(
&self,
name: &str,