diff --git a/crates/persisting-pchronicle-cli/src/lib.rs b/crates/persisting-pchronicle-cli/src/lib.rs index a03a7d154..df2949ad1 100644 --- a/crates/persisting-pchronicle-cli/src/lib.rs +++ b/crates/persisting-pchronicle-cli/src/lib.rs @@ -2244,7 +2244,8 @@ async fn run_serve( .with_context(|| format!("bind pChronicle Warehouse to {listen}"))?; let warehouse = if let Some(path) = args.catalog_config.as_ref() { let acl = server::catalog::CatalogAcl::load(path)?; - server::PreparedWarehouse::prepare_catalog(acl, config.clone()).await? + server::PreparedWarehouse::prepare_catalog(acl, config.clone(), Some(path.clone())) + .await? } else if args.gateway.is_some() { server::PreparedWarehouse::prepare_live(config.clone()).await? } else { diff --git a/crates/persisting-pchronicle-cli/src/server/catalog.rs b/crates/persisting-pchronicle-cli/src/server/catalog.rs index a9945e032..15aee69ab 100644 --- a/crates/persisting-pchronicle-cli/src/server/catalog.rs +++ b/crates/persisting-pchronicle-cli/src/server/catalog.rs @@ -1,5 +1,6 @@ use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet}; use std::path::Path; +use std::sync::Arc; use anyhow::{Context, Result, anyhow}; use persisting_pchronicle::storage::{DatasetLocation, DatasetMount}; @@ -8,6 +9,10 @@ use url::Url; use super::problem::ApiError; +#[path = "catalog_reload.rs"] +mod reload; +pub(crate) use reload::{CatalogSnapshot, CatalogState}; + pub(crate) const ACCESS_KEY_HEADER: &str = "x-pchronicle-access-key"; pub(crate) const SECRET_KEY_HEADER: &str = "x-pchronicle-secret-key"; const MAX_CATALOG_CONFIG_BYTES: u64 = 1024 * 1024; @@ -26,14 +31,14 @@ pub(crate) struct CatalogLibrary { pub secret_key: Option, } -#[derive(Debug, Clone)] +#[derive(Debug, Clone, PartialEq, Eq)] pub(crate) struct CatalogUser { pub name: String, secret_key: String, datasets: Vec, } -#[derive(Debug, Clone)] +#[derive(Debug, Clone, PartialEq, Eq)] pub(crate) struct CatalogAcl { libraries: BTreeMap, users_by_access_key: HashMap, @@ -151,6 +156,13 @@ impl CatalogAcl { .ok_or_else(catalog_unauthorized) } + pub(crate) fn public_mounts(&self) -> Vec { + self.public_for_all() + .into_iter() + .filter_map(|library| DatasetMount::new(&library.name, &library.uri).ok()) + .collect() + } + pub(crate) fn public_for_all(&self) -> Vec { self.public_datasets .iter() @@ -808,13 +820,12 @@ fn parent_handles_path(path: &str) -> bool { } pub(super) async fn list_datasets( - axum::extract::State(state): axum::extract::State, + snapshot: Option>>, headers: axum::http::HeaderMap, ) -> Result>, ApiError> { - let acl = state - .catalog_acl - .as_ref() - .ok_or_else(|| ApiError::not_found("catalog is not enabled"))?; + let axum::Extension(snapshot) = + snapshot.ok_or_else(|| ApiError::not_found("catalog is not enabled"))?; + let acl = &snapshot.acl; let has_credential_headers = headers.contains_key(ACCESS_KEY_HEADER) || headers.contains_key(SECRET_KEY_HEADER); let libraries = match credentials_from_headers(&headers) { @@ -829,14 +840,13 @@ pub(super) async fn list_datasets( } pub(super) async fn get_dataset( - axum::extract::State(state): axum::extract::State, + snapshot: Option>>, axum::extract::Path(name): axum::extract::Path, headers: axum::http::HeaderMap, ) -> Result, ApiError> { - let acl = state - .catalog_acl - .as_ref() - .ok_or_else(|| ApiError::not_found("catalog is not enabled"))?; + let axum::Extension(snapshot) = + snapshot.ok_or_else(|| ApiError::not_found("catalog is not enabled"))?; + let acl = &snapshot.acl; let user = acl.authenticate_headers(&headers)?; let ticket = acl .ticket_for(user, &name) @@ -851,11 +861,20 @@ pub(super) async fn catalog_data_plane_layer( ) -> axum::response::Response { use axum::response::IntoResponse; - if state.catalog_query_worker || state.catalog_acl.is_none() { + if state.catalog_query_worker { return next.run(request).await; } + let Some(catalog) = &state.catalog_acl else { + return next.run(request).await; + }; let path = request.uri().path().to_owned(); - if !path.starts_with("/api/") || parent_handles_path(&path) { + if !path.starts_with("/api/") { + return next.run(request).await; + } + let snapshot = catalog.snapshot(); + let acl = &snapshot.acl; + request.extensions_mut().insert(snapshot.clone()); + if parent_handles_path(&path) { return next.run(request).await; } // Anonymous browsing is limited to wildcard-granted datasets. @@ -864,12 +883,10 @@ pub(super) async fn catalog_data_plane_layer( && url::form_urlencoded::parse(request.uri().query().unwrap_or("").as_bytes()) .any(|(key, value)| key == "ui" && (value == "true" || value == "1")) { - if let Some(library) = state.catalog_acl.as_ref().and_then(|acl| { - acl.libraries.values().find(|library| { - acl.public_for_all() - .iter() - .any(|public| public.name == library.name) - }) + if let Some(library) = acl.libraries.values().find(|library| { + acl.public_for_all() + .iter() + .any(|public| public.name == library.name) }) { apply_library_env(library); } @@ -899,31 +916,27 @@ pub(super) async fn catalog_data_plane_layer( path, "catalog explorer request entering data-plane router" ); - let public = state.catalog_acl.as_ref().is_some_and(|acl| { + let public = { dataset.is_empty() || acl .public_for_all() .iter() .any(|library| library.name == dataset) - }); + }; if !public { - return dispatch_query_worker(&state, request) + return dispatch_query_worker(&state, acl, request) .await .unwrap_or_else(|error| error.into_response()); } if !dataset.is_empty() { - if let Some(library) = state.catalog_acl.as_ref().and_then(|acl| { - acl.libraries - .values() - .find(|library| library.name == dataset) - }) { + if let Some(library) = acl + .libraries + .values() + .find(|library| library.name == dataset) + { apply_library_env(library); } - if let Some((access_key, secret_key)) = state - .catalog_acl - .as_ref() - .and_then(|acl| acl.credentials_for_public(dataset)) - { + if let Some((access_key, secret_key)) = acl.credentials_for_public(dataset) { let headers = request.headers_mut(); if let (Ok(access_key), Ok(secret_key)) = (access_key.parse(), secret_key.parse()) { headers.insert(ACCESS_KEY_HEADER, access_key); @@ -937,17 +950,14 @@ pub(super) async fn catalog_data_plane_layer( ); return next.run(request).await; } - let mounts = state - .catalog_acl - .as_ref() - .unwrap() + let mounts = acl .public_for_all() .into_iter() .filter_map(|library| DatasetMount::new(&library.name, &library.uri).ok()) .collect::>(); return axum::Json(super::explorer::catalog_tree_from_mount_specs(&mounts)).into_response(); } - match dispatch_query_worker(&state, request).await { + match dispatch_query_worker(&state, acl, request).await { Ok(response) => response, Err(error) => error.into_response(), } @@ -955,13 +965,10 @@ pub(super) async fn catalog_data_plane_layer( async fn dispatch_query_worker( state: &super::AppState, + acl: &CatalogAcl, request: axum::http::Request, ) -> Result { use super::catalog_worker::{WorkerRequest, validate_backends}; - let acl = state - .catalog_acl - .as_ref() - .ok_or_else(|| ApiError::not_found("catalog is not enabled"))?; let (access_key, secret_key) = credentials_from_headers(request.headers()).ok_or_else(catalog_unauthorized)?; let user = acl @@ -1064,7 +1071,7 @@ async fn dispatch_query_worker( mod tests { use super::*; - const SAMPLE: &str = r#" + pub(super) const SAMPLE: &str = r#" [datasets.prod] uri = "s3://bucket/prod" endpoint = "http://127.0.0.1:9000" @@ -1362,9 +1369,10 @@ uri = "{}" let acl = CatalogAcl::load(&catalog).unwrap(); let config = crate::server::ChronicleServerConfig::mounted(acl.mounts().unwrap()).unwrap(); - let warehouse = crate::server::PreparedWarehouse::prepare_catalog(acl, config) - .await - .unwrap(); + let warehouse = + crate::server::PreparedWarehouse::prepare_catalog(acl, config, Some(catalog.clone())) + .await + .unwrap(); assert!(warehouse.dataset_names().is_empty()); assert!(warehouse.state.catalog.read().await.is_none()); use tower::ServiceExt; @@ -1393,12 +1401,14 @@ uri = "{}" let acl = CatalogAcl::load(&catalog).unwrap(); let config = crate::server::ChronicleServerConfig::front_only(); - let warehouse = crate::server::PreparedWarehouse::prepare_catalog(acl, config) - .await - .unwrap(); + let warehouse = + crate::server::PreparedWarehouse::prepare_catalog(acl, config, Some(catalog.clone())) + .await + .unwrap(); assert!(warehouse.state.config.datasets.is_empty()); - assert_eq!(warehouse.state.browse_mounts.len(), 1); - assert_eq!(warehouse.state.browse_mounts[0].name, "shared"); + let snapshot = warehouse.state.catalog_acl.as_ref().unwrap().snapshot(); + assert_eq!(snapshot.acl.public_mounts().len(), 1); + assert_eq!(snapshot.acl.public_mounts()[0].name, "shared"); } #[tokio::test] @@ -1426,9 +1436,10 @@ uri = "{}" let mut config = crate::server::ChronicleServerConfig::mounted(acl.mounts().unwrap()).unwrap(); config.home_links = vec![crate::server::parse_home_link("Realtime=/litefuse").unwrap()]; - let warehouse = crate::server::PreparedWarehouse::prepare_catalog(acl, config) - .await - .unwrap(); + let warehouse = + crate::server::PreparedWarehouse::prepare_catalog(acl, config, Some(catalog.clone())) + .await + .unwrap(); let response = warehouse .router() .oneshot( @@ -1567,6 +1578,270 @@ uri = "{}" assert_eq!(status, axum::http::StatusCode::NOT_FOUND); } + async fn reload_warehouse(path: &Path) -> crate::server::PreparedWarehouse { + crate::server::PreparedWarehouse::prepare_catalog( + CatalogAcl::load(path).unwrap(), + crate::server::ChronicleServerConfig::front_only(), + Some(path.to_path_buf()), + ) + .await + .unwrap() + } + + async fn wait_for_acl( + warehouse: &crate::server::PreparedWarehouse, + predicate: impl Fn(&CatalogAcl) -> bool, + ) { + tokio::time::timeout(std::time::Duration::from_secs(5), async { + loop { + let snapshot = warehouse.state.catalog_acl.as_ref().unwrap().snapshot(); + if predicate(&snapshot.acl) { + break; + } + tokio::time::sleep(std::time::Duration::from_millis(20)).await; + } + }) + .await + .expect("catalog reload did not become visible"); + } + + async fn ticket_status( + warehouse: &crate::server::PreparedWarehouse, + access: &str, + secret: &str, + ) -> axum::http::StatusCode { + use tower::ServiceExt; + warehouse + .router() + .oneshot(catalog_request( + "/api/v1/catalog/datasets/prod", + Some(access), + Some(secret), + )) + .await + .unwrap() + .status() + } + + #[tokio::test] + async fn catalog_acl_reload_picks_up_new_user_without_restarting() { + let temporary = tempfile::tempdir().unwrap(); + let path = temporary.path().join("catalog.toml"); + std::fs::write(&path, SAMPLE).unwrap(); + let warehouse = reload_warehouse(&path).await; + let user = issue_user(&path, "new_user").unwrap(); + grant_datasets(&path, "new_user", &["prod".into()]).unwrap(); + wait_for_acl(&warehouse, |acl| { + acl.authenticate(&user.access_key, &user.secret_key) + .is_some_and(|user| acl.ticket_for(user, "prod").is_some()) + }) + .await; + assert_eq!( + ticket_status(&warehouse, &user.access_key, &user.secret_key).await, + axum::http::StatusCode::OK + ); + assert_eq!( + ticket_status(&warehouse, "USER_AK", "USER_SK").await, + axum::http::StatusCode::OK + ); + + // A request that already acquired a snapshot can finish with its old grant. + let before_revoke = warehouse.state.catalog_acl.as_ref().unwrap().snapshot(); + revoke_datasets(&path, "new_user", &["prod".into()]).unwrap(); + wait_for_acl(&warehouse, |acl| { + acl.authenticate(&user.access_key, &user.secret_key) + .is_some_and(|user| acl.ticket_for(user, "prod").is_none()) + }) + .await; + assert_eq!( + ticket_status(&warehouse, &user.access_key, &user.secret_key).await, + axum::http::StatusCode::NOT_FOUND + ); + let old_user = before_revoke + .acl + .authenticate(&user.access_key, &user.secret_key) + .unwrap(); + assert!(before_revoke.acl.ticket_for(old_user, "prod").is_some()); + } + + #[tokio::test] + async fn catalog_acl_reload_keeps_inflight_request_snapshot() { + use tower::ServiceExt; + let temporary = tempfile::tempdir().unwrap(); + let path = temporary.path().join("catalog.toml"); + std::fs::write(&path, SAMPLE).unwrap(); + let warehouse = reload_warehouse(&path).await; + let entered = Arc::new(tokio::sync::Notify::new()); + let resume = Arc::new(tokio::sync::Notify::new()); + let app = axum::Router::new() + .route( + "/api/v1/catalog/datasets/prod", + axum::routing::get({ + let entered = entered.clone(); + let resume = resume.clone(); + move |axum::Extension(snapshot): axum::Extension>| { + let entered = entered.clone(); + let resume = resume.clone(); + async move { + entered.notify_one(); + resume.notified().await; + let user = snapshot.acl.authenticate("USER_AK", "USER_SK").unwrap(); + assert!(snapshot.acl.ticket_for(user, "prod").is_some()); + axum::http::StatusCode::OK + } + } + }), + ) + .layer(axum::middleware::from_fn_with_state( + warehouse.state.clone(), + catalog_data_plane_layer, + )); + let inflight = tokio::spawn(app.oneshot(catalog_request( + "/api/v1/catalog/datasets/prod", + Some("USER_AK"), + Some("USER_SK"), + ))); + tokio::time::timeout(std::time::Duration::from_secs(5), entered.notified()) + .await + .unwrap(); + revoke_datasets(&path, "alice", &["prod".into()]).unwrap(); + wait_for_acl(&warehouse, |acl| { + acl.ticket_for(acl.authenticate("USER_AK", "USER_SK").unwrap(), "prod") + .is_none() + }) + .await; + assert_eq!( + ticket_status(&warehouse, "USER_AK", "USER_SK").await, + axum::http::StatusCode::NOT_FOUND + ); + resume.notify_one(); + let response = tokio::time::timeout(std::time::Duration::from_secs(5), inflight) + .await + .unwrap() + .unwrap() + .unwrap(); + assert_eq!(response.status(), axum::http::StatusCode::OK); + } + + #[tokio::test] + async fn catalog_acl_reload_preserves_service_on_failure_and_recovers() { + use tower::ServiceExt; + let temporary = tempfile::tempdir().unwrap(); + let path = temporary.path().join("catalog.toml"); + std::fs::write(&path, SAMPLE).unwrap(); + let warehouse = reload_warehouse(&path).await; + // Exercise both invalid contents and a missing file across polling ticks. + for missing in [false, true] { + if missing { + std::fs::remove_file(&path).unwrap(); + } else { + std::fs::write(&path, "secret_key = \"SENSITIVE\" invalid").unwrap(); + } + tokio::time::sleep(std::time::Duration::from_millis(1100)).await; + assert_eq!( + ticket_status(&warehouse, "USER_AK", "USER_SK").await, + axum::http::StatusCode::OK + ); + for uri in ["/", "/api/health", "/api/ui"] { + assert_eq!( + warehouse + .router() + .oneshot(catalog_request(uri, None, None)) + .await + .unwrap() + .status(), + axum::http::StatusCode::OK + ); + } + } + std::fs::write(&path, SAMPLE.replace("USER_SK", "ROTATED_SK")).unwrap(); + wait_for_acl(&warehouse, |acl| { + acl.authenticate("USER_AK", "ROTATED_SK").is_some() + }) + .await; + assert_eq!( + ticket_status(&warehouse, "USER_AK", "USER_SK").await, + axum::http::StatusCode::UNAUTHORIZED + ); + assert_eq!( + ticket_status(&warehouse, "USER_AK", "ROTATED_SK").await, + axum::http::StatusCode::OK + ); + } + + #[tokio::test] + async fn catalog_acl_reload_rejects_backend_changes_as_a_whole() { + let temporary = tempfile::tempdir().unwrap(); + let path = temporary.path().join("catalog.toml"); + std::fs::write(&path, SAMPLE).unwrap(); + let warehouse = reload_warehouse(&path).await; + std::fs::write( + &path, + SAMPLE + .replace("USER_SK", "ROTATED_SK") + .replace("s3://bucket/prod", "s3://bucket/other"), + ) + .unwrap(); + tokio::time::sleep(std::time::Duration::from_millis(1100)).await; + assert_eq!( + ticket_status(&warehouse, "USER_AK", "USER_SK").await, + axum::http::StatusCode::OK + ); + assert_eq!( + ticket_status(&warehouse, "USER_AK", "ROTATED_SK").await, + axum::http::StatusCode::UNAUTHORIZED + ); + std::fs::write(&path, SAMPLE.replace("USER_SK", "ROTATED_SK")).unwrap(); + wait_for_acl(&warehouse, |acl| { + acl.authenticate("USER_AK", "ROTATED_SK").is_some() + }) + .await; + } + + #[tokio::test] + async fn catalog_acl_reload_updates_public_browse_snapshot() { + use tower::ServiceExt; + let temporary = tempfile::tempdir().unwrap(); + let path = temporary.path().join("catalog.toml"); + let dataset = temporary.path().join("data"); + std::fs::create_dir_all(dataset.join("visible_child")).unwrap(); + let private = format!("[datasets.prod]\nuri = \"{}\"\n", dataset.display()); + std::fs::write(&path, &private).unwrap(); + let warehouse = reload_warehouse(&path).await; + std::fs::write( + &path, + format!("{private}\n[[grants]]\nuser = \"*\"\ndataset = \"prod\"\n"), + ) + .unwrap(); + wait_for_acl(&warehouse, |acl| !acl.public_for_all().is_empty()).await; + let (status, body) = catalog_body( + warehouse + .router() + .oneshot(catalog_request( + "/api/explorer/tree?dataset=prod", + None, + None, + )) + .await + .unwrap(), + ) + .await; + assert_eq!(status, axum::http::StatusCode::OK); + assert!(body.contains("visible_child"), "{body}"); + std::fs::write(&path, &private).unwrap(); + wait_for_acl(&warehouse, |acl| acl.public_for_all().is_empty()).await; + let response = warehouse + .router() + .oneshot(catalog_request( + "/api/explorer/tree?dataset=prod", + None, + None, + )) + .await + .unwrap(); + assert_eq!(response.status(), axum::http::StatusCode::UNAUTHORIZED); + } + #[tokio::test] async fn catalog_data_plane_requires_headers_without_spawning_worker() { use tower::ServiceExt; diff --git a/crates/persisting-pchronicle-cli/src/server/catalog_reload.rs b/crates/persisting-pchronicle-cli/src/server/catalog_reload.rs new file mode 100644 index 000000000..7d35ba018 --- /dev/null +++ b/crates/persisting-pchronicle-cli/src/server/catalog_reload.rs @@ -0,0 +1,187 @@ +use std::io::Read; +use std::path::{Path, PathBuf}; +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use tokio::sync::{OnceCell, watch}; + +use super::{CatalogAcl, MAX_CATALOG_CONFIG_BYTES}; +use crate::server::ui_cache::BrowseCoordinator; + +const RELOAD_INTERVAL: Duration = Duration::from_secs(3); +const ERROR_LOG_INTERVAL: Duration = Duration::from_secs(60); + +/// Authentication and public browsing use the same immutable request snapshot. +pub(crate) struct CatalogSnapshot { + pub(crate) acl: CatalogAcl, + browse: Arc>, +} + +impl CatalogSnapshot { + pub(crate) async fn browse(&self) -> &BrowseCoordinator { + self.browse + .get_or_init(|| BrowseCoordinator::start(self.acl.public_mounts())) + .await + } +} + +pub(crate) struct CatalogState { + current: watch::Receiver>, + task: Option>, +} + +impl Drop for CatalogState { + fn drop(&mut self) { + if let Some(task) = &self.task { + task.abort(); + } + } +} + +impl CatalogState { + pub(crate) fn new(acl: CatalogAcl, path: Option) -> Self { + let snapshot = Arc::new(CatalogSnapshot { + acl, + browse: Arc::new(OnceCell::new()), + }); + let (sender, current) = watch::channel(snapshot); + let task = path.map(|path| tokio::spawn(reload_loop(path, sender))); + Self { current, task } + } + + pub(crate) fn snapshot(&self) -> Arc { + self.current.borrow().clone() + } +} + +// Return only fixed diagnostic categories: TOML errors can include source lines +// containing credentials, and validation errors can include user-supplied text. +fn read_update( + path: &Path, + previous: Option, +) -> Result, &'static str> { + let file = std::fs::File::open(path).map_err(|_| "read_failed")?; + let metadata = file.metadata().map_err(|_| "read_failed")?; + if !metadata.is_file() || metadata.len() > MAX_CATALOG_CONFIG_BYTES { + return Err("invalid_file"); + } + let mut content = String::new(); + file.take(MAX_CATALOG_CONFIG_BYTES + 1) + .read_to_string(&mut content) + .map_err(|_| "read_failed")?; + if content.len() as u64 > MAX_CATALOG_CONFIG_BYTES { + return Err("invalid_file"); + } + let hash = blake3::hash(content.as_bytes()); + if previous == Some(hash) { + return Ok(None); + } + let document = super::parse_catalog_file(&content).map_err(|_| "invalid_toml")?; + let acl = CatalogAcl::from_document(document).map_err(|_| "invalid_catalog")?; + Ok(Some((hash, acl))) +} + +async fn reload_loop(path: PathBuf, sender: watch::Sender>) { + let mut interval = tokio::time::interval(RELOAD_INTERVAL); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + let mut fingerprint = None; + let mut last_error_log = None; + loop { + interval.tick().await; + let path = path.clone(); + let result = tokio::task::spawn_blocking(move || read_update(&path, fingerprint)) + .await + .unwrap_or(Err("reload_task_failed")); + let result = match result { + Ok(Some((hash, acl))) => { + let previous = sender.borrow().clone(); + if acl.libraries != previous.acl.libraries { + Err("dataset_changes_require_restart") + } else { + fingerprint = Some(hash); + if acl != previous.acl { + let browse = if acl.public_datasets == previous.acl.public_datasets { + previous.browse.clone() + } else { + Arc::new(OnceCell::new()) + }; + sender.send_replace(Arc::new(CatalogSnapshot { acl, browse })); + tracing::info!(target: "pchronicle.serve", "catalog ACL reloaded"); + } + Ok(()) + } + } + Ok(None) => Ok(()), + Err(reason) => Err(reason), + }; + match result { + Ok(()) => { + if last_error_log.take().is_some() { + tracing::info!(target: "pchronicle.serve", "catalog config reload recovered"); + } + } + Err(reason) => { + if last_error_log.is_none_or(|last: Instant| last.elapsed() >= ERROR_LOG_INTERVAL) { + tracing::error!( + target: "pchronicle.serve", + reason, + "catalog config reload rejected; keeping last valid ACL; requested grants and revocations have not taken effect" + ); + last_error_log = Some(Instant::now()); + } + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn reload_errors_never_include_config_values() { + let temporary = tempfile::tempdir().unwrap(); + let path = temporary.path().join("catalog.toml"); + std::fs::write(&path, "secret_key = \"SECRET_MUST_NOT_BE_LOGGED\" invalid").unwrap(); + assert_eq!(read_update(&path, None).unwrap_err(), "invalid_toml"); + std::fs::write(&path, "[datasets.prod]\nuri = \"/tmp/data\"\n[users.SECRET_MUST_NOT_BE_LOGGED]\naccess_key = \"\"\nsecret_key = \"secret\"\n").unwrap(); + assert_eq!(read_update(&path, None).unwrap_err(), "invalid_catalog"); + std::fs::remove_file(&path).unwrap(); + assert_eq!(read_update(&path, None).unwrap_err(), "read_failed"); + } + + #[test] + fn reload_detects_atomic_replacement_and_bounds_input() { + let temporary = tempfile::tempdir().unwrap(); + let path = temporary.path().join("catalog.toml"); + std::fs::write(&path, super::super::tests::SAMPLE).unwrap(); + let (hash, _) = read_update(&path, None).unwrap().unwrap(); + assert!(read_update(&path, Some(hash)).unwrap().is_none()); + let replacement = path.with_extension("tmp"); + std::fs::write( + &replacement, + super::super::tests::SAMPLE.replace("USER_SK", "NEXT_SK"), + ) + .unwrap(); + std::fs::rename(replacement, &path).unwrap(); + let (_, acl) = read_update(&path, Some(hash)).unwrap().unwrap(); + assert!(acl.authenticate("USER_AK", "NEXT_SK").is_some()); + std::fs::File::create(&path) + .unwrap() + .set_len(MAX_CATALOG_CONFIG_BYTES + 1) + .unwrap(); + assert_eq!(read_update(&path, Some(hash)).unwrap_err(), "invalid_file"); + } + + #[tokio::test] + async fn dropping_catalog_stops_reload_task() { + let temporary = tempfile::tempdir().unwrap(); + let path = temporary.path().join("catalog.toml"); + std::fs::write(&path, super::super::tests::SAMPLE).unwrap(); + let state = CatalogState::new(CatalogAcl::load(&path).unwrap(), Some(path)); + let task = state.task.as_ref().unwrap().abort_handle(); + drop(state); + tokio::task::yield_now().await; + assert!(task.is_finished()); + } +} diff --git a/crates/persisting-pchronicle-cli/src/server/mod.rs b/crates/persisting-pchronicle-cli/src/server/mod.rs index 1c3676663..157701369 100644 --- a/crates/persisting-pchronicle-cli/src/server/mod.rs +++ b/crates/persisting-pchronicle-cli/src/server/mod.rs @@ -63,7 +63,7 @@ struct AppState { /// Gateway-backed Warehouses read canonical events from the latest /// manifest for single-trace observation, independent of projection idle. live_reads: bool, - catalog_acl: Option>, + catalog_acl: Option>, catalog_query_worker: bool, catalog_workers: Arc, browse_mounts: Arc>, @@ -253,7 +253,7 @@ impl PreparedWarehouse { #[cfg(test)] pub(crate) async fn prepare_catalog_front(acl: catalog::CatalogAcl) -> anyhow::Result { - Self::prepare_catalog(acl, ChronicleServerConfig::front_only()).await + Self::prepare_catalog(acl, ChronicleServerConfig::front_only(), None).await } /// The listener authenticates and brokers credentials; only workers mount @@ -261,21 +261,14 @@ impl PreparedWarehouse { pub(crate) async fn prepare_catalog( acl: catalog::CatalogAcl, mut config: ChronicleServerConfig, + catalog_config: Option, ) -> anyhow::Result { - let browse_mounts = acl - .public_for_all() - .into_iter() - .filter_map(|library| DatasetMount::new(&library.name, &library.uri).ok()) - .collect::>(); config.datasets.clear(); config.default_dataset = None; let mut state = app_state(config); - state.browse_mounts = Arc::new(browse_mounts); - state.catalog_acl = Some(Arc::new(acl)); - // The front process owns the public browse cache. Query workers remain - // isolated, but UI navigation must be available without mounting data - // in the front process. - browse_coordinator(&state).await; + let catalog = Arc::new(catalog::CatalogState::new(acl, catalog_config)); + catalog.snapshot().browse().await; + state.catalog_acl = Some(catalog); Ok(Self { state }) } @@ -1364,6 +1357,7 @@ fn preview_needle(query: &str) -> String { async fn explorer_tree( State(state): State, + snapshot: Option>>, request_id: RequestId, metrics: RequestMetrics, query: Result, QueryRejection>, @@ -1391,18 +1385,24 @@ async fn explorer_tree( metrics.record("browse", started); return Ok(Json(serde_json::to_value(view).unwrap())); }; - let owned_mount; - let mount = if let Some(mount) = state.browse_mounts.iter().find(|mount| mount.name == name) { - mount - } else if let Some(library) = state.catalog_acl.as_ref().and_then(|acl| { - acl.public_for_all() + if let Some(axum::Extension(snapshot)) = snapshot { + let mount = snapshot + .acl + .public_mounts() .into_iter() - .find(|library| library.name == name) - }) { - owned_mount = DatasetMount::new(&library.name, &library.uri) + .find(|mount| mount.name == name) + .ok_or_else(|| ApiError::not_found("dataset not found"))?; + let started = Instant::now(); + let view = snapshot + .browse() + .await + .tree(&mount, prefix) + .await .map_err(|error| fail(&request_id, "explorer_tree", error))?; - &owned_mount - } else { + metrics.record("browse", started); + return Ok(Json(serde_json::to_value(view).unwrap())); + } + let Some(mount) = state.browse_mounts.iter().find(|mount| mount.name == name) else { return Ok(Json( serde_json::to_value(explorer::catalog_tree_from_path_list(name, prefix, &[])).unwrap(), )); diff --git a/docs/src/en/rfcs/0013-pchronicle-warehouse-catalog.md b/docs/src/en/rfcs/0013-pchronicle-warehouse-catalog.md index b989b33e2..4cde65186 100644 --- a/docs/src/en/rfcs/0013-pchronicle-warehouse-catalog.md +++ b/docs/src/en/rfcs/0013-pchronicle-warehouse-catalog.md @@ -159,7 +159,7 @@ datasets = ["evals"] ## CLI 签发与授权 -签发和改授权是 **写 `catalog.toml` 的 CLI**,不是运行中 Warehouse 的 HTTP API。出现 `catalog` 子命令时 MUST NOT 启动 listener。正在运行的 serve MUST 重启后才读到新用户或新授权。 +签发和改授权是 **写 `catalog.toml` 的 CLI**,不是运行中 Warehouse 的 HTTP API。出现 `catalog` 子命令时 MUST NOT 启动 listener。运行中的 serve 每 3 秒检查配置,用户和授权无需重启即可生效。 ```text pchronicle serve catalog dataset add --catalog-config FILE NAME --uri URI [--endpoint URL] [--region REGION] [--access-key KEY] [--secret-key KEY] @@ -198,7 +198,7 @@ pchronicle serve --catalog-config FILE --listen 127.0.0.1:8081 - `revoke` 从该用户的 `datasets` 里去掉列出的名字。未知用户、或该用户当前并未持有的 library 名 MUST 失败。 - 两个命令的 stdout 只报 `name` 与更新后的 `datasets`,MUST NOT 打印密钥。 -改写配置可以整表重写,不要求保留注释。新用户在重启 serve 之前无法登录。 +改写配置可以整表重写,不要求保留注释。新用户通常在 3 秒内生效。 ## HTTP diff --git a/docs/src/zh/pchronicle/guides/serve.md b/docs/src/zh/pchronicle/guides/serve.md index ead680d54..52ef6ec76 100644 --- a/docs/src/zh/pchronicle/guides/serve.md +++ b/docs/src/zh/pchronicle/guides/serve.md @@ -64,7 +64,7 @@ pchronicle serve --catalog-config catalog.toml --listen 127.0.0.1:8081 `catalog.toml` 列出 libraries(`[datasets.*]`,本地 path 或 `s3://`)和 users。 `serve catalog dataset add|remove|list` 改写 libraries,不启动 HTTP。 `serve catalog issue` 写入一个无授权用户,并把 sk 只打印到这次 stdout; -`grant` / `revoke` 改该用户可打开的 library 名称。改文件后必须重启 serve。 +`grant` / `revoke` 改该用户可打开的 library 名称。运行中的 serve 每 3 秒检查配置,无需重启。 NAME 使用 `*` 可以为当前所有用户授予或撤销数据集。新建用户不会自动继承过去的通配授权,创建后请重新执行命令。 `pchronicle serve --catalog-config` 会把文件中的 **全部** library 挂进 Warehouse diff --git a/docs/src/zh/rfcs/0013-pchronicle-warehouse-catalog.md b/docs/src/zh/rfcs/0013-pchronicle-warehouse-catalog.md index 0bbdef4a8..f7807c4ad 100644 --- a/docs/src/zh/rfcs/0013-pchronicle-warehouse-catalog.md +++ b/docs/src/zh/rfcs/0013-pchronicle-warehouse-catalog.md @@ -62,7 +62,7 @@ pchronicle query @team/prod 'SELECT 1' ### 非目标 - STS、临时凭证轮换、或把用户钥映射成短时 AWS session。 -- 热加载 `catalog.toml`;改配置 MUST 重启 serve。 +- 热加载 Dataset 定义和 S3 后端凭证;这类变更仍需重启 serve。 - 在运行中的 Warehouse 上提供 HTTP 签发接口。 - 提供独立 `catalog serve` 二进制。 - 在已运行的 Tokio runtime 上 `fork(2)`(未定义行为)。