Skip to content

Commit f7da3c0

Browse files
authored
fix: add a new middleware for peer communication (#1558)
Adds a prism_endpoint CLI option and Prism URL wiring, introduces IntraClusterRequest middleware and per-user intra-cluster header propagation, threads userid through many sync/forward flows, renames route params (username→userid, stream_name→logstream), and adds RBAC/user helper and visibility changes.
1 parent f1278d2 commit f7da3c0

22 files changed

Lines changed: 524 additions & 180 deletions

src/alerts/alert_types.rs

Lines changed: 23 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818

1919
use std::{str::FromStr, time::Duration};
2020

21-
use actix_web::http::header::{HeaderMap, HeaderName, HeaderValue};
21+
use actix_web::http::header::{HeaderName, HeaderValue};
2222
use chrono::{DateTime, Utc};
2323
use serde_json::Value;
2424
use tonic::async_trait;
@@ -36,9 +36,12 @@ use crate::{
3636
get_number_of_agg_exprs,
3737
target::{self, NotificationConfig},
3838
},
39-
handlers::http::query::create_streams_for_distributed,
39+
handlers::http::{
40+
middleware::{CLUSTER_SECRET, CLUSTER_SECRET_HEADER},
41+
query::create_streams_for_distributed,
42+
},
4043
metastore::metastore_traits::MetastoreObject,
41-
parseable::PARSEABLE,
44+
parseable::{DEFAULT_TENANT, PARSEABLE},
4245
query::resolve_stream_names,
4346
rbac::map::SessionKey,
4447
storage::object_storage::alert_json_path,
@@ -87,18 +90,29 @@ impl MetastoreObject for ThresholdAlert {
8790
impl AlertTrait for ThresholdAlert {
8891
async fn eval_alert(&self) -> Result<Option<String>, AlertError> {
8992
let time_range = extract_time_range(&self.eval_config)?;
90-
let auth = if let Some(tenant) = self.tenant_id.as_ref()
91-
&& let Some(header) = TENANT_METADATA.get_global_query_auth(tenant)
92-
{
93-
let mut map = HeaderMap::new();
93+
94+
let tenant = self.tenant_id.as_deref().unwrap_or(DEFAULT_TENANT);
95+
let auth = if let Some((_, hash)) = CLUSTER_SECRET.get() {
96+
let mut map = actix_web::http::header::HeaderMap::new();
97+
if let Some(header) = TENANT_METADATA.get_global_query_auth(tenant) {
98+
map.insert(
99+
HeaderName::from_static("authorization"),
100+
HeaderValue::from_str(&header).unwrap(),
101+
);
102+
}
103+
map.insert(
104+
HeaderName::from_static(CLUSTER_SECRET_HEADER),
105+
HeaderValue::from_str(hash).unwrap(),
106+
);
94107
map.insert(
95-
HeaderName::from_static("authorization"),
96-
HeaderValue::from_str(&header).unwrap(),
108+
HeaderName::from_static("intra-cluster-tenant"),
109+
HeaderValue::from_str(tenant).unwrap(),
97110
);
98111
Some(map)
99112
} else {
100113
None
101114
};
115+
102116
let query_result =
103117
execute_alert_query(auth, self.get_query(), &time_range, &self.tenant_id).await?;
104118

src/cli.rs

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -407,6 +407,14 @@ pub struct Options {
407407
)]
408408
pub querier_endpoint: String,
409409

410+
#[arg(
411+
long,
412+
env = "P_PRISM_ENDPOINT",
413+
default_value = "",
414+
help = "URL to connect to the prism node. Default is the address of the server"
415+
)]
416+
pub prism_endpoint: String,
417+
410418
#[command(flatten)]
411419
pub oidc: Option<OidcConfig>,
412420

@@ -593,6 +601,7 @@ impl Options {
593601
Mode::Ingest => self.get_endpoint(&self.ingestor_endpoint, "P_INGESTOR_ENDPOINT"),
594602
Mode::Index => self.get_endpoint(&self.indexer_endpoint, "P_INDEXER_ENDPOINT"),
595603
Mode::Query => self.get_endpoint(&self.querier_endpoint, "P_QUERIER_ENDPOINT"),
604+
Mode::Prism => self.get_endpoint(&self.prism_endpoint, "P_PRISM_ENDPOINT"),
596605
_ => return self.build_url(&self.address),
597606
};
598607

src/handlers/http/cluster/mod.rs

Lines changed: 81 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,8 @@
1919
pub mod utils;
2020
use actix_web::http::StatusCode;
2121
use actix_web::http::header::HeaderMap;
22+
use base64::Engine;
23+
use base64::prelude::BASE64_STANDARD;
2224
use futures::{StreamExt, future, stream};
2325
use http::header;
2426
use lazy_static::lazy_static;
@@ -50,7 +52,9 @@ use crate::rbac::role::model::Role;
5052
use crate::rbac::user::User;
5153
use crate::stats::Stats;
5254
use crate::storage::{ObjectStorageError, ObjectStoreFormat};
53-
use crate::utils::{create_intracluster_auth_headermap, get_tenant_id_from_request};
55+
use crate::utils::{
56+
create_intracluster_auth_headermap, get_tenant_id_from_request, get_user_from_request,
57+
};
5458

5559
use super::base_path_without_preceding_slash;
5660
use super::ingest::PostError;
@@ -519,6 +523,7 @@ pub async fn sync_users_with_roles_with_ingestors(
519523
let userid = userid.to_owned();
520524
let headers = req.headers().clone();
521525
let op = operation.to_string();
526+
let caller_userid = get_user_from_request(req).unwrap();
522527
for_each_live_node(tenant_id, move |ingestor| {
523528
let url = format!(
524529
"{}{}/user/{}/role/sync/{}",
@@ -529,7 +534,8 @@ pub async fn sync_users_with_roles_with_ingestors(
529534
);
530535

531536
let role_data = role_data.clone();
532-
let headermap = create_intracluster_auth_headermap(&headers, &ingestor.token);
537+
let headermap =
538+
create_intracluster_auth_headermap(&headers, &ingestor.token, &caller_userid);
533539
async move {
534540
let res = INTRA_CLUSTER_CLIENT
535541
.patch(url)
@@ -568,6 +574,7 @@ pub async fn sync_user_deletion_with_ingestors(
568574
tenant_id: &Option<String>,
569575
) -> Result<(), RBACError> {
570576
let userid = userid.to_owned();
577+
let caller_userid = get_user_from_request(req).unwrap();
571578
let headers = req.headers().clone();
572579
for_each_live_node(tenant_id, move |ingestor| {
573580
let url = format!(
@@ -576,7 +583,8 @@ pub async fn sync_user_deletion_with_ingestors(
576583
base_path_without_preceding_slash(),
577584
userid
578585
);
579-
let headermap = create_intracluster_auth_headermap(&headers, &ingestor.token);
586+
let headermap =
587+
create_intracluster_auth_headermap(&headers, &ingestor.token, &caller_userid);
580588
async move {
581589
let res = INTRA_CLUSTER_CLIENT
582590
.delete(url)
@@ -625,6 +633,7 @@ pub async fn sync_user_creation(
625633
RBACError::SerdeError(err)
626634
})?;
627635

636+
let caller_userid = get_user_from_request(req)?;
628637
let userid = userid.to_string();
629638
let headers = req.headers().clone();
630639
for_each_live_node(tenant_id, move |node| {
@@ -634,7 +643,7 @@ pub async fn sync_user_creation(
634643
base_path_without_preceding_slash(),
635644
userid
636645
);
637-
let headermap = create_intracluster_auth_headermap(&headers, &node.token);
646+
let headermap = create_intracluster_auth_headermap(&headers, &node.token, &caller_userid);
638647
let user_data = user_data.clone();
639648

640649
async move {
@@ -672,20 +681,23 @@ pub async fn sync_password_reset_with_ingestors(
672681
req: HttpRequest,
673682
username: &str,
674683
) -> Result<(), RBACError> {
675-
let username = username.to_owned();
684+
let userid = username.to_owned();
676685
let tenant_id = get_tenant_id_from_request(&req);
686+
let caller_userid = get_user_from_request(&req).unwrap();
687+
let headers = req.headers().clone();
677688
for_each_live_node(&tenant_id, move |ingestor| {
678689
let url = format!(
679690
"{}{}/user/{}/generate-new-password/sync",
680691
ingestor.domain_name,
681692
base_path_without_preceding_slash(),
682-
username
693+
userid
683694
);
684-
695+
let headermap =
696+
create_intracluster_auth_headermap(&headers, &ingestor.token, &caller_userid);
685697
async move {
686698
let res = INTRA_CLUSTER_CLIENT
687699
.post(url)
688-
.header(header::AUTHORIZATION, &ingestor.token)
700+
.headers(headermap)
689701
.header(header::CONTENT_TYPE, "application/json")
690702
.send()
691703
.await
@@ -713,11 +725,14 @@ pub async fn sync_password_reset_with_ingestors(
713725

714726
// forward the put role request to all ingestors and queriers to keep them in sync
715727
pub async fn sync_role_update(
728+
req: &HttpRequest,
716729
name: String,
717730
role: Role,
718731
tenant_id: &Option<String>,
719732
) -> Result<(), RoleError> {
720733
let tenant = tenant_id.to_owned();
734+
let userid = get_user_from_request(req).unwrap();
735+
let headers = req.headers().clone();
721736
for_each_live_node(tenant_id, move |node| {
722737
let url = format!(
723738
"{}{}/role/{}/sync",
@@ -727,12 +742,12 @@ pub async fn sync_role_update(
727742
);
728743

729744
let role = role.clone();
730-
745+
let headermap = create_intracluster_auth_headermap(&headers, &node.token, &userid);
731746
let tenant_id = tenant.clone();
732747
async move {
733748
let res = INTRA_CLUSTER_CLIENT
734749
.put(url)
735-
.header(header::AUTHORIZATION, &node.token)
750+
.headers(headermap)
736751
.header(header::CONTENT_TYPE, "application/json")
737752
.json(&SyncRole::new(role, tenant_id.clone()))
738753
.send()
@@ -759,6 +774,52 @@ pub async fn sync_role_update(
759774
.await
760775
}
761776

777+
// forward the put role request to all ingestors and queriers to keep them in sync
778+
pub async fn sync_role_delete(
779+
req: &HttpRequest,
780+
name: String,
781+
tenant_id: &Option<String>,
782+
) -> Result<(), RoleError> {
783+
let userid = get_user_from_request(req).unwrap();
784+
let headers = req.headers().clone();
785+
for_each_live_node(tenant_id, move |node| {
786+
let url = format!(
787+
"{}{}/role/{}/sync",
788+
node.domain_name,
789+
base_path_without_preceding_slash(),
790+
name
791+
);
792+
793+
let headermap = create_intracluster_auth_headermap(&headers, &node.token, &userid);
794+
async move {
795+
let res = INTRA_CLUSTER_CLIENT
796+
.delete(url)
797+
.headers(headermap)
798+
.header(header::CONTENT_TYPE, "application/json")
799+
.send()
800+
.await
801+
.map_err(|err| {
802+
error!(
803+
"Fatal: failed to forward request to node: {}\n Error: {:?}",
804+
node.domain_name, err
805+
);
806+
RoleError::Network(err)
807+
})?;
808+
809+
if !res.status().is_success() {
810+
error!(
811+
"failed to forward request to node: {}\nResponse Returned: {:?}",
812+
node.domain_name,
813+
res.text().await
814+
);
815+
}
816+
817+
Ok(())
818+
}
819+
})
820+
.await
821+
}
822+
762823
pub fn fetch_daily_stats(
763824
date: &str,
764825
stream_meta_list: &[ObjectStoreFormat],
@@ -1904,7 +1965,6 @@ pub async fn send_query_request(
19041965
let mut map = reqwest::header::HeaderMap::new();
19051966

19061967
if let Some(auth) = auth_token {
1907-
// always basic auth
19081968
for (key, value) in auth.iter() {
19091969
if let Ok(name) = reqwest::header::HeaderName::from_bytes(key.as_str().as_bytes())
19101970
&& let Ok(val) = reqwest::header::HeaderValue::from_bytes(value.as_bytes())
@@ -1918,6 +1978,16 @@ pub async fn send_query_request(
19181978
reqwest::header::HeaderValue::from_str(&querier.token).unwrap(),
19191979
);
19201980
};
1981+
if map.get("intra-cluster-userid").is_none() {
1982+
let token: Vec<&str> = querier.token().split(' ').collect();
1983+
let decode = BASE64_STANDARD.decode(token[1]).unwrap();
1984+
let user = String::from_utf8(decode).unwrap();
1985+
let user = user.split_once(':').unwrap();
1986+
map.insert(
1987+
reqwest::header::HeaderName::from_static("intra-cluster-userid"),
1988+
reqwest::header::HeaderValue::from_str(user.0).unwrap(),
1989+
);
1990+
}
19211991
let res = match INTRA_CLUSTER_CLIENT
19221992
.post(uri)
19231993
.timeout(Duration::from_secs(300))

src/handlers/http/ingest.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -434,10 +434,10 @@ pub async fn handle_otel_traces_ingestion(
434434
// fails if the logstream does not exist
435435
pub async fn post_event(
436436
req: HttpRequest,
437-
stream_name: Path<String>,
437+
logstream: Path<String>,
438438
Json(json): Json<StrictValue>,
439439
) -> Result<HttpResponse, PostError> {
440-
let stream_name = stream_name.into_inner();
440+
let stream_name = logstream.into_inner();
441441
let tenant_id = get_tenant_id_from_request(&req);
442442
let internal_stream_names = PARSEABLE.streams.list_internal_streams(&tenant_id);
443443
if internal_stream_names.contains(&stream_name) {

0 commit comments

Comments
 (0)