diff --git a/crates/persisting-pchronicle-cli/src/exchange/import.rs b/crates/persisting-pchronicle-cli/src/exchange/import.rs index a2144df14..554c3cef5 100644 --- a/crates/persisting-pchronicle-cli/src/exchange/import.rs +++ b/crates/persisting-pchronicle-cli/src/exchange/import.rs @@ -587,6 +587,12 @@ pub(crate) async fn run_import( progress.notice(&line)?; } progress.flush_log(stderr)?; + if let Some(wal) = &wal + && let Ok(guard) = wal.lock() + && let Err(error) = guard.remove() + { + tracing::warn!(target: "pchronicle.import", %error, "failed to remove completed import WAL"); + } Ok(()) } @@ -1938,9 +1944,20 @@ pub(crate) async fn finalize_storyline_import_indexes( store: &StorylineLanceStore, progress: &mut CliProgress, ) -> Result<()> { - progress - .stage(StageId::Commit) - .set_current("optimize indices (final)"); + let commit = progress.stage(StageId::Commit); + commit.set_current("data committed; building indices"); + let heartbeat = commit.clone(); + let started = std::time::Instant::now(); + let heartbeat_task = tokio::spawn(async move { + let mut tick = tokio::time::interval(std::time::Duration::from_secs(2)); + loop { + tick.tick().await; + heartbeat.set_activity_override(&format!( + "building indices · elapsed={}s", + started.elapsed().as_secs() + )); + } + }); let _index_progress = progress.attach_index_progress(); store .maintain(&persisting_pchronicle::storage::LanceMaintenanceOptions { @@ -1951,9 +1968,9 @@ pub(crate) async fn finalize_storyline_import_indexes( }) .await .context("finalize Storyline indexes after progressive import")?; - progress - .stage(StageId::Commit) - .set_current("optimize indices done"); + heartbeat_task.abort(); + commit.set_activity_override("indices built"); + commit.set_current("indices built"); Ok(()) } @@ -1968,7 +1985,20 @@ pub(crate) async fn commit_storyline_import_batch( anyhow::ensure!(!batch.is_empty(), "storyline import commit batch is empty"); let batch_len = batch.len() as u64; commit.set_current(format!("batch={batch_len}")); - let report = match append_generation.as_deref() { + let heartbeat = commit.clone(); + let heartbeat_started = std::time::Instant::now(); + let heartbeat_task = tokio::spawn(async move { + let mut tick = tokio::time::interval(std::time::Duration::from_secs(2)); + loop { + tick.tick().await; + heartbeat.set_activity_override(&format!( + "writing Lance snapshot · batch={batch_len} · elapsed={}s", + heartbeat_started.elapsed().as_secs() + )); + } + }); + let report_result: Result<_> = async { + Ok(match append_generation.as_deref() { Some(generation) => { tracing::info!( committed_before = committed_storylines, @@ -2010,7 +2040,11 @@ pub(crate) async fn commit_storyline_import_batch( ) })? } - }; + }) + } + .await; + heartbeat_task.abort(); + let report = report_result?; anyhow::ensure!( report.storylines as u64 == batch_len, "storyline import batch report does not match batch size" diff --git a/crates/persisting-pchronicle-cli/src/exchange/wal.rs b/crates/persisting-pchronicle-cli/src/exchange/wal.rs index 472f2f23e..4e773283e 100644 --- a/crates/persisting-pchronicle-cli/src/exchange/wal.rs +++ b/crates/persisting-pchronicle-cli/src/exchange/wal.rs @@ -160,6 +160,15 @@ impl ImportWal { self.failed.len() } + /// A completed import no longer needs recovery state. + pub(crate) fn remove(&self) -> Result<()> { + if self.dir.exists() { + fs::remove_dir_all(&self.dir) + .with_context(|| format!("remove completed import WAL {}", self.dir.display()))?; + } + Ok(()) + } + pub(crate) fn skip_paths(&self) -> HashSet { self.done .iter() @@ -365,5 +374,8 @@ mod tests { .unwrap(); assert!(!reset.should_skip("a.json")); assert_eq!(reset.done_count(), 0); + let dir = reset.dir().to_path_buf(); + reset.remove().unwrap(); + assert!(!dir.exists()); } } diff --git a/crates/persisting-pchronicle-cli/src/lib.rs b/crates/persisting-pchronicle-cli/src/lib.rs index 261c3f80e..7a59e71ff 100644 --- a/crates/persisting-pchronicle-cli/src/lib.rs +++ b/crates/persisting-pchronicle-cli/src/lib.rs @@ -117,29 +117,19 @@ impl Cli { } } -/// Apply S3 backend keys from `--catalog-config` and local `@name` pin settings -/// before the multi-threaded Tokio runtime starts. `std::env::set_var` after -/// worker threads exist is racy on macOS and can leave OpenDAL unable to see -/// `AWS_REGION`. -pub fn apply_catalog_backend_env_before_runtime(cli: &Cli) -> Result<()> { - apply_serve_catalog_backend_env(cli)?; - apply_command_pin_backend_env(cli)?; - Ok(()) +/// Run the internal exec worker before constructing a runtime. Ordinary CLI +/// invocations return None and retain their existing execution path. +pub fn run_catalog_worker_before_runtime(cli: &Cli) -> Option> { + match &cli.command { + Command::Serve(args) if args.catalog_query_worker => Some(server::catalog_worker::run()), + _ => None, + } } -fn apply_serve_catalog_backend_env(cli: &Cli) -> Result<()> { - let Command::Serve(args) = &cli.command else { - return Ok(()); - }; - if args.command.is_some() || args.catalog_query_worker { - return Ok(()); - } - let Some(path) = args.catalog_config.as_ref() else { - return Ok(()); - }; - let acl = server::catalog::CatalogAcl::load(path)?; - acl.apply_backend_env(); - Ok(()) +/// Apply local @name pin credentials before creating runtime threads. Catalog +/// server credentials are passed exclusively to isolated workers through IPC. +pub fn apply_catalog_backend_env_before_runtime(cli: &Cli) -> Result<()> { + apply_command_pin_backend_env(cli) } fn apply_command_pin_backend_env(cli: &Cli) -> Result<()> { @@ -1030,18 +1020,16 @@ struct ServeArgs { #[arg(long = "gateway-debug", alias = "debug", requires = "gateway_config")] debug: bool, - /// Directory ACL file (libraries + users). Mounts every [datasets.*] entry - /// into Warehouse and enables catalog:// locators. Mutually exclusive with - /// positional Dataset mounts. Apply S3 endpoint/region/keys from the file - /// before opening stores. + /// Directory ACL file. Authenticate API requests and execute them in + /// bounded, user-scoped worker processes with explicit backend credentials. #[arg( long = "catalog-config", value_name = "FILE", - conflicts_with_all = ["config", "storage", "positional_storage"] + conflicts_with_all = ["config", "storage", "positional_storage", "gateway", "gateway_config", "gateway_dataset", "control"] )] catalog_config: Option, - /// Internal: run one filtered Warehouse request from stdin and exit. + /// Internal: serve framed Warehouse requests over private stdin/stdout IPC. #[arg(long = "catalog-query-worker", hide = true)] catalog_query_worker: bool, } @@ -1644,7 +1632,7 @@ pub async fn run_with_stdio( }) => run_echo(args, &mut diagnostics).await, Command::Serve(args) => { if args.catalog_query_worker { - return server::catalog::run_catalog_query_worker().await; + bail!("catalog worker must start before the async runtime"); } if let Some(command) = args.command { return run_serve_catalog(command, stdout_is_terminal, stdout, &mut diagnostics); @@ -2454,9 +2442,12 @@ fn resolve_serve_config_with_settings( let storage = serve_storage_uris(args); let gateway_dataset = resolve_gateway_dataset_uri(args, settings_override)?; let mut config = if let Some(path) = args.catalog_config.as_deref() { - let acl = server::catalog::CatalogAcl::load(path)?; - acl.apply_backend_env(); - server::ChronicleServerConfig::mounted(acl.mounts()?)? + anyhow::ensure!( + gateway_dataset.is_none() && args.control.is_none(), + "catalog workers cannot share Gateway or Control listeners" + ); + server::catalog::CatalogAcl::load(path)?; + server::ChronicleServerConfig::front_only() } else { match (args.config.as_deref(), storage.as_slice()) { (Some(config), []) => { diff --git a/crates/persisting-pchronicle-cli/src/main.rs b/crates/persisting-pchronicle-cli/src/main.rs index cc21a4846..019dcbbcd 100644 --- a/crates/persisting-pchronicle-cli/src/main.rs +++ b/crates/persisting-pchronicle-cli/src/main.rs @@ -10,6 +10,13 @@ use persisting_pchronicle_cli::{ fn main() -> ExitCode { let cli = Cli::parse(); + if let Some(result) = persisting_pchronicle_cli::run_catalog_worker_before_runtime(&cli) { + return if result.is_ok() { + ExitCode::SUCCESS + } else { + ExitCode::from(1) + }; + } let debug_errors = cli.debug_errors(); // OpenDAL/Lance read AWS_* from the process environment. Applying catalog // backend keys after the multi-threaded Tokio runtime starts is racy on diff --git a/crates/persisting-pchronicle-cli/src/server/catalog.rs b/crates/persisting-pchronicle-cli/src/server/catalog.rs index 414c85471..a9945e032 100644 --- a/crates/persisting-pchronicle-cli/src/server/catalog.rs +++ b/crates/persisting-pchronicle-cli/src/server/catalog.rs @@ -37,6 +37,7 @@ pub(crate) struct CatalogUser { pub(crate) struct CatalogAcl { libraries: BTreeMap, users_by_access_key: HashMap, + public_datasets: BTreeSet, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -111,12 +112,20 @@ impl CatalogAcl { fn from_document(file: CatalogFile) -> Result { let libraries = build_libraries(&file)?; + let public_datasets = file + .grants + .iter() + .filter(|grant| grant.user == "*") + .map(|grant| grant.dataset.clone()) + .collect(); Ok(Self { users_by_access_key: build_users(&file, &libraries)?, libraries, + public_datasets, }) } + #[cfg(test)] pub(crate) fn mounts(&self) -> Result> { self.libraries .values() @@ -124,16 +133,6 @@ impl CatalogAcl { .collect() } - pub(crate) fn apply_backend_env(&self) { - if let Some(library) = self - .libraries - .values() - .find(|library| library.access_key.is_some()) - { - apply_library_env(library); - } - } - pub(crate) fn authenticate(&self, access_key: &str, secret_key: &str) -> Option<&CatalogUser> { let user = self.users_by_access_key.get(access_key)?; if !secret_keys_match(&user.secret_key, secret_key) { @@ -152,6 +151,25 @@ impl CatalogAcl { .ok_or_else(catalog_unauthorized) } + pub(crate) fn public_for_all(&self) -> Vec { + self.public_datasets + .iter() + .filter_map(|name| self.libraries.get(name)) + .map(CatalogLibraryPublic::from) + .collect() + } + + fn credentials_for_public(&self, dataset: &str) -> Option<(&str, &str)> { + self.users_by_access_key + .iter() + .find_map(|(access_key, user)| { + user.datasets + .iter() + .any(|name| name == dataset) + .then_some((access_key.as_str(), user.secret_key.as_str())) + }) + } + pub(crate) fn list_for(&self, user: &CatalogUser) -> Vec { user.datasets .iter() @@ -198,59 +216,78 @@ pub(crate) fn issue_user(path: &Path, name: &str) -> Result { } pub(crate) fn grant_datasets(path: &Path, name: &str, datasets: &[String]) -> Result> { - let name = canonical_user_name(name)?; let mut file = load_editable_catalog(path)?; let library_names = canonical_library_names(&file)?; - anyhow::ensure!(file.users.contains_key(&name), "unknown user '{name}'"); + let users: Vec = if name == "*" { + file.users.keys().cloned().collect() + } else { + vec![canonical_user_name(name)?] + }; + for user in &users { + anyhow::ensure!(file.users.contains_key(user), "unknown user '{user}'"); + } for dataset in datasets { let dataset = granted_library_name(&library_names, dataset)?; - if !file - .grants - .iter() - .any(|grant| grant.user == name && grant.dataset == dataset) - { - file.grants.push(CatalogGrantFile { - user: name.clone(), - dataset, - permissions: vec!["read".into(), "query".into(), "analyze".into()], - }); + for user in &users { + if !file + .grants + .iter() + .any(|grant| grant.user == *user && grant.dataset == dataset) + { + file.grants.push(CatalogGrantFile { + user: user.clone(), + dataset: dataset.clone(), + permissions: vec!["read".into(), "query".into(), "analyze".into()], + }); + } } } + let mut seen = HashSet::new(); let granted = file .grants .iter() - .filter(|grant| grant.user == name) + .filter(|grant| users.iter().any(|user| user == &grant.user)) .map(|grant| grant.dataset.clone()) + .filter(|dataset| seen.insert(dataset.clone())) .collect(); write_catalog_file(path, &file)?; Ok(granted) } pub(crate) fn revoke_datasets(path: &Path, name: &str, datasets: &[String]) -> Result> { - let name = canonical_user_name(name)?; let mut file = load_editable_catalog(path)?; - anyhow::ensure!(file.users.contains_key(&name), "unknown user '{name}'"); - let mut to_remove = Vec::new(); - for dataset in datasets { - let dataset = DatasetMount::new(dataset, "validation") - .with_context(|| format!("catalog library name '{dataset}'"))? - .name; + let users: Vec = if name == "*" { + file.users.keys().cloned().collect() + } else { + vec![canonical_user_name(name)?] + }; + for user in &users { + anyhow::ensure!(file.users.contains_key(user), "unknown user '{user}'"); + } + let to_remove = datasets + .iter() + .map(|dataset| DatasetMount::new(dataset, "validation").map(|mount| mount.name)) + .collect::>>()?; + for dataset in &to_remove { anyhow::ensure!( file.grants .iter() - .any(|grant| grant.user == name && grant.dataset == dataset), + .any(|grant| users.iter().any(|user| user == &grant.user) + && grant.dataset == *dataset), "catalog user '{name}' does not grant '{dataset}'" ); - to_remove.push(dataset); } file.grants.retain(|grant| { - !(grant.user == name && to_remove.iter().any(|dataset| dataset == &grant.dataset)) + !(users.iter().any(|user| user == &grant.user) + && to_remove.iter().any(|dataset| dataset == &grant.dataset)) }); + let mut seen = HashSet::new(); let remaining = file .grants .iter() - .filter(|grant| grant.user == name) + .filter(|grant| users.iter().any(|user| user == &grant.user)) .map(|grant| grant.dataset.clone()) + .filter(|dataset| seen.insert(dataset.clone())) .collect(); write_catalog_file(path, &file)?; Ok(remaining) @@ -515,24 +552,31 @@ fn build_users( let mut users_by_access_key = HashMap::new(); let mut datasets_by_user: HashMap> = HashMap::new(); for grant in &file.grants { - anyhow::ensure!( - file.users.contains_key(&grant.user), - "catalog grant references unknown user '{}'", - grant.user - ); - anyhow::ensure!( - libraries.contains_key(&grant.dataset), - "catalog grant references unknown dataset '{}'", - grant.dataset - ); - let entry = datasets_by_user.entry(grant.user.clone()).or_default(); - anyhow::ensure!( - !entry.contains(&grant.dataset), - "catalog grant for user '{}' and dataset '{}' is duplicated", - grant.user, - grant.dataset - ); - entry.push(grant.dataset.clone()); + let grant_users: Vec<&str> = if grant.user == "*" { + file.users.keys().map(String::as_str).collect() + } else { + vec![grant.user.as_str()] + }; + for grant_user in grant_users { + anyhow::ensure!( + file.users.contains_key(grant_user), + "catalog grant references unknown user '{}'", + grant_user + ); + anyhow::ensure!( + libraries.contains_key(&grant.dataset), + "catalog grant references unknown dataset '{}'", + grant.dataset + ); + let entry = datasets_by_user.entry(grant_user.to_owned()).or_default(); + anyhow::ensure!( + !entry.contains(&grant.dataset), + "catalog grant for user '{}' and dataset '{}' is duplicated", + grant_user, + grant.dataset + ); + entry.push(grant.dataset.clone()); + } } for (name, user) in &file.users { let access_key = user.access_key.trim().to_owned(); @@ -771,8 +815,17 @@ pub(super) async fn list_datasets( .catalog_acl .as_ref() .ok_or_else(|| ApiError::not_found("catalog is not enabled"))?; - let user = acl.authenticate_headers(&headers)?; - Ok(axum::Json(acl.list_for(user))) + let has_credential_headers = + headers.contains_key(ACCESS_KEY_HEADER) || headers.contains_key(SECRET_KEY_HEADER); + let libraries = match credentials_from_headers(&headers) { + Some((access, secret)) => acl + .authenticate(&access, &secret) + .map(|user| acl.list_for(user)) + .ok_or_else(catalog_unauthorized)?, + None if !has_credential_headers => acl.public_for_all(), + None => return Err(catalog_unauthorized()), + }; + Ok(axum::Json(libraries)) } pub(super) async fn get_dataset( @@ -793,213 +846,218 @@ pub(super) async fn get_dataset( pub(super) async fn catalog_data_plane_layer( axum::extract::State(state): axum::extract::State, - request: axum::http::Request, + mut request: axum::http::Request, next: axum::middleware::Next, ) -> axum::response::Response { use axum::response::IntoResponse; - if state.catalog_query_worker - || state.catalog_acl.is_none() - || !state.config.datasets.is_empty() - { + if state.catalog_query_worker || state.catalog_acl.is_none() { return next.run(request).await; } let path = request.uri().path().to_owned(); if !path.starts_with("/api/") || parent_handles_path(&path) { return next.run(request).await; } + // Anonymous browsing is limited to wildcard-granted datasets. + if credentials_from_headers(request.headers()).is_none() + && path.ends_with("/query/tables") + && 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) + }) + }) { + apply_library_env(library); + } + let catalog = super::QueryCatalog { + snapshot_id: String::new(), + read_only: true, + database: String::new(), + storage_path: String::new(), + path_column: "_file_", + datasets: Vec::new(), + tables: super::query_table_summaries(), + }; + return axum::Json(catalog).into_response(); + } + if credentials_from_headers(request.headers()).is_none() && path.ends_with("/explorer/tree") { + let params: BTreeMap<_, _> = + url::form_urlencoded::parse(request.uri().query().unwrap_or("").as_bytes()) + .into_owned() + .collect(); + let dataset = params + .get("dataset") + .map(String::as_str) + .unwrap_or_default(); + tracing::info!( + target: "pchronicle.serve", + dataset, + path, + "catalog explorer request entering data-plane router" + ); + let public = state.catalog_acl.as_ref().is_some_and(|acl| { + dataset.is_empty() + || acl + .public_for_all() + .iter() + .any(|library| library.name == dataset) + }); + if !public { + return dispatch_query_worker(&state, 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) + }) { + apply_library_env(library); + } + if let Some((access_key, secret_key)) = state + .catalog_acl + .as_ref() + .and_then(|acl| 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); + headers.insert(SECRET_KEY_HEADER, secret_key); + } + } + tracing::info!( + target: "pchronicle.serve", + dataset, + "serving public catalog explorer request from front browse cache" + ); + return next.run(request).await; + } + let mounts = state + .catalog_acl + .as_ref() + .unwrap() + .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 { Ok(response) => response, Err(error) => error.into_response(), } } -#[derive(Debug, Serialize, Deserialize)] -struct CatalogWorkerJob { - request_id: String, - method: String, - path: String, - query: String, - body: Vec, - mounts: Vec, -} - -#[derive(Debug, Serialize, Deserialize)] -struct CatalogWorkerResult { - status: u16, - #[serde(default)] - content_type: Option, - body: Vec, -} - async fn dispatch_query_worker( state: &super::AppState, 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 user = acl.authenticate_headers(request.headers())?.clone(); - tracing::debug!( - target: super::problem::LOG_TARGET, - user = %user.name, - libraries = user.datasets.len(), - path = %request.uri().path(), - "dispatch catalog query worker" - ); - if user.datasets.is_empty() { + let (access_key, secret_key) = + credentials_from_headers(request.headers()).ok_or_else(catalog_unauthorized)?; + let user = acl + .authenticate(&access_key, &secret_key) + .ok_or_else(catalog_unauthorized)?; + let _permit = state.catalog_workers.admit()?; + let selected = url::form_urlencoded::parse(request.uri().query().unwrap_or("").as_bytes()) + .filter(|(name, _)| name == "dataset") + .map(|(_, value)| value.into_owned()) + .collect::>(); + if selected.len() > 1 { + return Err(ApiError::invalid_request("duplicate dataset parameter")); + } + let selected = selected.first().filter(|name| !name.is_empty()); + if let Some(name) = selected + && acl.ticket_for(user, name).is_none() + { return Err(ApiError::not_found("dataset not found")); } - let mounts: Vec = user + let mut mounts: Vec = user .datasets .iter() - .filter_map(|name| acl.ticket_for(&user, name)) + .filter_map(|name| acl.ticket_for(user, name)) .collect(); - let request_id = request - .extensions() - .get::() - .map(|id| id.0.clone()) - .unwrap_or_default(); - let method = request.method().as_str().to_owned(); - let path = request.uri().path().to_owned(); - let query = request.uri().query().unwrap_or("").to_owned(); - let body = axum::body::to_bytes(request.into_body(), 1024 * 1024) - .await - .map_err(|error| ApiError::invalid_request(format!("read catalog query body: {error}")))?; - let job = CatalogWorkerJob { - request_id, - method, - path, - query, - body: body.to_vec(), - mounts, - }; - let payload = serde_json::to_vec(&job) - .map_err(|error| ApiError::internal("", "catalog_worker", anyhow::anyhow!(error)))?; - let exe = std::env::current_exe() - .map_err(|error| ApiError::internal("", "catalog_worker", anyhow::anyhow!(error)))?; - let mut child = tokio::process::Command::new(exe) - .arg("serve") - .arg("--catalog-query-worker") - .stdin(std::process::Stdio::piped()) - .stdout(std::process::Stdio::piped()) - .stderr(std::process::Stdio::piped()) - .spawn() - .map_err(|error| ApiError::internal("", "catalog_worker", anyhow::anyhow!(error)))?; - { - use tokio::io::AsyncWriteExt; - let mut stdin = child.stdin.take().ok_or_else(|| { - ApiError::internal( - "", - "catalog_worker", - anyhow::anyhow!("catalog query worker stdin is missing"), - ) - })?; - stdin - .write_all(&payload) - .await - .map_err(|error| ApiError::internal("", "catalog_worker", anyhow::anyhow!(error)))?; - stdin - .shutdown() - .await - .map_err(|error| ApiError::internal("", "catalog_worker", anyhow::anyhow!(error)))?; + if mounts.is_empty() { + return Err(ApiError::not_found("dataset not found")); } - let output = tokio::time::timeout(std::time::Duration::from_secs(60), child.wait_with_output()) - .await - .map_err(|_| ApiError::unavailable())? - .map_err(|error| ApiError::internal("", "catalog_worker", anyhow::anyhow!(error)))?; - if !output.status.success() { - let stderr = String::from_utf8_lossy(&output.stderr); - tracing::error!( - target: super::problem::LOG_TARGET, - handler = "catalog_worker", - exit = output.status.code().unwrap_or(-1), - stderr = %super::problem::truncate_utf8(&stderr, super::problem::QUERY_LOG_LIMIT), - "warehouse request failed" - ); - return Err(ApiError::internal( - "", - "catalog_worker", - anyhow::anyhow!("catalog query worker exited unsuccessfully"), - )); - } - let result: CatalogWorkerResult = serde_json::from_slice(&output.stdout) - .map_err(|error| ApiError::internal("", "catalog_worker", anyhow::anyhow!(error)))?; - let status = axum::http::StatusCode::from_u16(result.status) - .unwrap_or(axum::http::StatusCode::INTERNAL_SERVER_ERROR); - let mut response = axum::response::Response::new(axum::body::Body::from(result.body)); - *response.status_mut() = status; - if let Some(content_type) = result.content_type - && let Ok(value) = axum::http::HeaderValue::from_str(&content_type) + // Keep one stable worker for all of this user's compatible mounts. Only + // split scopes when the backend's process-global credentials require it. + if validate_backends(&mounts).is_err() + && let Some(name) = selected { - response - .headers_mut() - .insert(axum::http::header::CONTENT_TYPE, value); - } - Ok(response) -} - -pub(crate) async fn run_catalog_query_worker() -> Result<()> { - use std::io::{Read, Write}; - - use axum::body::Body; - use tower::ServiceExt; - - let mut stdin = Vec::new(); - std::io::stdin() - .read_to_end(&mut stdin) - .context("read catalog query worker job")?; - let job: CatalogWorkerJob = - serde_json::from_slice(&stdin).context("decode catalog query worker job")?; - if let Some(library) = job.mounts.first() { - apply_library_env(library); + mounts.retain(|mount| &mount.name == name); + } + validate_backends(&mounts).map_err(|error| ApiError::invalid_request(error.to_string()))?; + // Workers use a private working directory; retain the server's interpretation + // of relative local mounts without opening any store in the parent. + for mount in &mut mounts { + if let Some(path) = DatasetLocation::parse(&mount.uri) + .ok() + .and_then(|l| l.local_path().map(std::path::Path::to_path_buf)) + { + mount.uri = std::path::absolute(path) + .map_err(|error| ApiError::internal("", "catalog_worker", error.into()))? + .to_string_lossy() + .into_owned(); + } } - anyhow::ensure!(!job.mounts.is_empty(), "catalog query worker needs mounts"); - let mounts = job - .mounts - .iter() - .map(|library| DatasetMount::new(&library.name, &library.uri)) - .collect::>>()?; - let config = super::ChronicleServerConfig::mounted(mounts)?; - let warehouse = super::PreparedWarehouse::prepare_query_worker(config).await?; - let mut uri = job.path; - if !job.query.is_empty() { - uri.push('?'); - uri.push_str(&job.query); - } - let mut builder = axum::http::Request::builder() - .method(job.method.as_str()) - .uri(uri); - if !job.body.is_empty() { - builder = builder.header(axum::http::header::CONTENT_TYPE, "application/json"); - } - let request = builder - .body(Body::from(job.body)) - .context("build catalog query worker request")?; - let response = warehouse - .router() - .oneshot(request) - .await - .map_err(|error| anyhow!(error))?; - let status = response.status().as_u16(); - let content_type = response + // Never persist credentials in filenames. Include grants and credential + // epochs so revocation/rotation cannot reopen an older user's cache. + let identity = serde_json::to_vec(&( + &access_key, + &user.name, + &user.secret_key, + &user.datasets, + &mounts, + )) + .map_err(|error| ApiError::internal("", "catalog_worker", error.into()))?; + let scope = blake3::hash(&identity).to_hex().to_string(); + let mut headers = request .headers() - .get(axum::http::header::CONTENT_TYPE) - .and_then(|value| value.to_str().ok()) - .map(str::to_owned); - let body = axum::body::to_bytes(response.into_body(), 8 * 1024 * 1024) + .iter() + .filter(|(name, _)| matches!(name.as_str(), "content-type" | "accept")) + .filter_map(|(name, value)| { + value + .to_str() + .ok() + .map(|v| (name.to_string(), v.to_owned())) + }) + .collect::>(); + if let Some(id) = request.extensions().get::() { + headers.push(("x-request-id".into(), id.0.clone())); + } + let method = request.method().to_string(); + let uri = request.uri().to_string(); + let body = tokio::time::timeout( + std::time::Duration::from_secs(10), + axum::body::to_bytes(request.into_body(), 1024 * 1024), + ) + .await + .map_err(|_| ApiError::unavailable())? + .map_err(|_| ApiError::invalid_request("catalog request body exceeds 1 MiB"))? + .to_vec(); + state + .catalog_workers + .execute( + scope, + mounts, + WorkerRequest { + method, + uri, + headers, + body, + }, + ) .await - .context("read catalog query worker response")?; - let result = CatalogWorkerResult { - status, - content_type, - body: body.to_vec(), - }; - serde_json::to_writer(std::io::stdout(), &result) - .context("write catalog query worker result")?; - std::io::stdout().flush().ok(); - Ok(()) } #[cfg(test)] @@ -1045,6 +1103,34 @@ dataset = "evals" permissions = ["read", "query"] "#; + #[test] + fn wildcard_grant_expands_to_all_users() { + let acl = CatalogAcl::parse( + r#" +[datasets.prod] +uri = "/tmp/prod" +[datasets.private] +uri = "/tmp/private" +[users.alice] +access_key = "a" +secret_key = "as" +[users.bob] +access_key = "b" +secret_key = "bs" +[[grants]] +user = "*" +dataset = "prod" +[[grants]] +user = "alice" +dataset = "private" +"#, + ) + .unwrap(); + assert_eq!(acl.list_for(acl.authenticate("a", "as").unwrap()).len(), 2); + assert_eq!(acl.list_for(acl.authenticate("b", "bs").unwrap()).len(), 1); + assert_eq!(acl.public_for_all().len(), 1); + } + #[test] fn parse_rejects_unknown_grant() { let error = CatalogAcl::parse( @@ -1251,7 +1337,7 @@ dataset = "prod" } #[tokio::test] - async fn prepare_catalog_mounts_local_datasets_without_credentials() { + async fn prepare_catalog_does_not_mount_datasets_in_parent() { let temporary = tempfile::tempdir().unwrap(); let left = temporary.path().join("left"); let right = temporary.path().join("right"); @@ -1279,7 +1365,40 @@ uri = "{}" let warehouse = crate::server::PreparedWarehouse::prepare_catalog(acl, config) .await .unwrap(); - assert_eq!(warehouse.dataset_names(), vec!["left", "right"]); + assert!(warehouse.dataset_names().is_empty()); + assert!(warehouse.state.catalog.read().await.is_none()); + use tower::ServiceExt; + let response = warehouse + .router() + .oneshot(catalog_request("/api/explorer/runs", None, None)) + .await + .unwrap(); + assert_eq!(response.status(), axum::http::StatusCode::UNAUTHORIZED); + } + + #[tokio::test] + async fn prepare_catalog_keeps_public_mounts_for_browse_cache() { + let temporary = tempfile::tempdir().unwrap(); + let dataset = temporary.path().join("shared"); + std::fs::create_dir_all(&dataset).unwrap(); + let catalog = temporary.path().join("catalog.toml"); + std::fs::write( + &catalog, + format!( + "[datasets.shared]\nuri = \"{}\"\n\n[[grants]]\nuser = \"*\"\ndataset = \"shared\"\n", + dataset.display() + ), + ) + .unwrap(); + + let acl = CatalogAcl::load(&catalog).unwrap(); + let config = crate::server::ChronicleServerConfig::front_only(); + let warehouse = crate::server::PreparedWarehouse::prepare_catalog(acl, config) + .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"); } #[tokio::test] @@ -1356,7 +1475,7 @@ uri = "{}" } #[tokio::test] - async fn catalog_list_requires_credentials_and_omits_backend_secrets() { + async fn catalog_list_allows_anonymous_public_datasets_and_rejects_invalid_credentials() { use tower::ServiceExt; let app = catalog_front().await; @@ -1367,7 +1486,7 @@ uri = "{}" .unwrap(), ) .await; - assert_eq!(status, axum::http::StatusCode::UNAUTHORIZED); + assert_eq!(status, axum::http::StatusCode::OK); let (status, body) = catalog_body( app.clone() @@ -1513,7 +1632,7 @@ secret_key = "BACKEND_SK" } #[test] - fn apply_catalog_backend_env_before_runtime_loads_s3_region() { + fn catalog_parent_does_not_install_backend_credentials() { use clap::Parser; let temporary = tempfile::tempdir().unwrap(); @@ -1540,13 +1659,20 @@ secret_key = "123" &catalog_arg, ]) .unwrap(); - unsafe { - std::env::remove_var("AWS_REGION"); - std::env::remove_var("AWS_DEFAULT_REGION"); - } + let keys = [ + "AWS_REGION", + "AWS_DEFAULT_REGION", + "AWS_ACCESS_KEY_ID", + "AWS_SECRET_ACCESS_KEY", + "AWS_ENDPOINT_URL_S3", + ]; + let before: Vec<_> = keys.iter().map(std::env::var_os).collect(); crate::apply_catalog_backend_env_before_runtime(&cli).unwrap(); - assert_eq!(std::env::var("AWS_REGION").unwrap(), "us-east-1"); - assert_eq!(std::env::var("AWS_ACCESS_KEY_ID").unwrap(), "123"); + let after: Vec<_> = keys.iter().map(std::env::var_os).collect(); + assert!( + before == after, + "parent must not mutate storage credentials" + ); } #[test] diff --git a/crates/persisting-pchronicle-cli/src/server/catalog_worker.rs b/crates/persisting-pchronicle-cli/src/server/catalog_worker.rs new file mode 100644 index 000000000..52343a515 --- /dev/null +++ b/crates/persisting-pchronicle-cli/src/server/catalog_worker.rs @@ -0,0 +1,500 @@ +//! Exec workers have an immutable authenticated scope. No storage client or +//! runtime is inherited from the listening process. +use std::{ + collections::HashMap, + io::{Read, Write}, + path::PathBuf, + process::Stdio, + sync::Arc, + time::Duration, +}; + +use anyhow::{Context, Result}; +use serde::{Deserialize, Serialize, de::DeserializeOwned}; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::sync::{Mutex, Semaphore}; +use tower::ServiceExt; + +use super::{ + catalog::{CatalogLibrary, apply_library_env}, + problem::ApiError, +}; + +const MAX_WORKERS: usize = 8; +const MAX_REQUESTS: usize = 32; +const FRAME_LIMIT: usize = 40 * 1024 * 1024; +const REQUEST_TIMEOUT: Duration = Duration::from_secs(60); + +#[derive(Serialize, Deserialize)] +pub(super) struct WorkerRequest { + pub method: String, + pub uri: String, + pub headers: Vec<(String, String)>, + pub body: Vec, +} + +#[derive(Serialize, Deserialize)] +pub(super) struct WorkerResponse { + status: u16, + headers: Vec<(String, String)>, + body: Vec, +} + +impl WorkerResponse { + pub(super) fn into_response(self) -> Result { + let mut response = axum::response::Response::new(axum::body::Body::from(self.body)); + *response.status_mut() = axum::http::StatusCode::from_u16(self.status)?; + for (name, value) in self.headers { + response + .headers_mut() + .append(name.parse::()?, value.parse()?); + } + Ok(response) + } +} + +#[derive(Serialize, Deserialize)] +struct Bootstrap { + mounts: Vec, +} + +type Slot = Arc>>; + +pub(super) struct WorkerPool { + slots: Mutex>, + requests: Semaphore, +} + +impl Default for WorkerPool { + fn default() -> Self { + Self { + slots: Mutex::new(HashMap::new()), + requests: Semaphore::new(MAX_REQUESTS), + } + } +} + +impl WorkerPool { + pub(super) fn admit(&self) -> Result, ApiError> { + self.requests + .try_acquire() + .map_err(|_| ApiError::unavailable()) + } + + async fn slot(&self, scope: &str) -> Result { + let mut slots = self.slots.lock().await; + if let Some(slot) = slots.get(scope) { + return Ok(slot.clone()); + } + if slots.len() >= MAX_WORKERS { + // Only evict a worker with no in-flight or queued request. Wait for + // its exit before spawning a replacement, keeping the process cap. + let idle = slots + .iter() + .find(|(_, slot)| Arc::strong_count(slot) == 1) + .map(|(key, _)| key.clone()); + let Some(idle) = idle else { + return Err(ApiError::unavailable()); + }; + if let Some(slot) = slots.remove(&idle) + && let Some(mut worker) = slot.lock().await.take() + { + let _ = worker.child.kill().await; + } + } + let slot = Arc::new(Mutex::new(None)); + slots.insert(scope.to_owned(), slot.clone()); + Ok(slot) + } + + pub(super) async fn execute( + &self, + scope: String, + mounts: Vec, + request: WorkerRequest, + ) -> Result { + tokio::time::timeout(REQUEST_TIMEOUT, async { + let slot = self.slot(&scope).await?; + let mut guard = slot.lock().await; + // Ownership stays in this future during IPC: cancellation, timeout + // or a partial frame drops/kills it rather than reusing dirty pipes. + let mut worker = match guard.take() { + Some(worker) => worker, + None => Worker::start(&scope, mounts).await.map_err(worker_error)?, + }; + let response = worker.exchange(&request).await.map_err(worker_error)?; + let response = response.into_response().map_err(worker_error)?; + *guard = Some(worker); + Ok(response) + }) + .await + .map_err(|_| ApiError::unavailable())? + } +} + +fn worker_error(error: anyhow::Error) -> ApiError { + // Protocol/OS diagnostics only; never log bootstrap payloads or child stderr. + ApiError::internal("", "catalog_worker", error) +} + +struct Worker { + child: tokio::process::Child, + input: tokio::process::ChildStdin, + output: tokio::process::ChildStdout, + _home: tempfile::TempDir, +} + +fn command( + exe: PathBuf, + home: &std::path::Path, + cache: &std::path::Path, +) -> tokio::process::Command { + let mut command = tokio::process::Command::new(exe); + command + .arg("serve") + .arg("--catalog-query-worker") + .env_clear() + .current_dir(home) + .env("HOME", home) + .env("USERPROFILE", home) + .env("XDG_CONFIG_HOME", home) + .env("XDG_CACHE_HOME", home) + .env("PCHRONICLE_CACHE_DIR", cache) + .env("AWS_EC2_METADATA_DISABLED", "true") + .env("RAYON_NUM_THREADS", "2") + .env("AWS_CONFIG_FILE", home.join("no-aws-config")) + .env( + "AWS_SHARED_CREDENTIALS_FILE", + home.join("no-aws-credentials"), + ) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + // Keep worker diagnostics visible to the serve process. Catalog UI + // requests run in this child, so hiding stderr also hid cache refresh + // failures and made a stalled browse impossible to diagnose. + .stderr(Stdio::inherit()) + .kill_on_drop(true); + // Retain only platform essentials and explicitly configured trust roots. + for key in [ + "PATH", + "SystemRoot", + "SSL_CERT_FILE", + "SSL_CERT_DIR", + "PCHRONICLE_QUERY_MEMORY_LIMIT", + "RUST_LOG", + ] { + if let Some(value) = std::env::var_os(key) { + command.env(key, value); + } + } + command +} + +impl Worker { + async fn start(scope: &str, mounts: Vec) -> Result { + let home = tempfile::tempdir()?; + let root = std::env::var_os("PCHRONICLE_CACHE_DIR") + .map(PathBuf::from) + .or_else(|| dirs::cache_dir().map(|p| p.join("pchronicle"))) + .context("no catalog worker cache directory")?; + let cache = std::path::absolute(root)?.join("workers").join(scope); + let mut builder = std::fs::DirBuilder::new(); + builder.recursive(true); + #[cfg(unix)] + { + use std::os::unix::fs::DirBuilderExt; + builder.mode(0o700); + } + builder.create(&cache)?; + let mut child = command(std::env::current_exe()?, home.path(), &cache).spawn()?; + let input = child.stdin.take().context("worker stdin missing")?; + let output = child.stdout.take().context("worker stdout missing")?; + let mut worker = Self { + child, + input, + output, + _home: home, + }; + worker.send(&Bootstrap { mounts }).await?; + let ready: bool = worker.receive().await?; + anyhow::ensure!(ready, "catalog worker failed to initialize"); + Ok(worker) + } + + async fn send(&mut self, message: &impl Serialize) -> Result<()> { + let bytes = encode(message)?; + self.input.write_all(&bytes).await?; + self.input.flush().await?; + Ok(()) + } + + async fn receive(&mut self) -> Result { + let size = self.output.read_u32().await? as usize; + anyhow::ensure!(size <= FRAME_LIMIT, "worker frame too large"); + let mut bytes = vec![0; size]; + self.output.read_exact(&mut bytes).await?; + Ok(serde_json::from_slice(&bytes)?) + } + + async fn exchange(&mut self, request: &WorkerRequest) -> Result { + self.send(request).await?; + self.receive().await + } +} + +fn encode(message: &impl Serialize) -> Result> { + let payload = serde_json::to_vec(message)?; + anyhow::ensure!(payload.len() <= FRAME_LIMIT, "worker frame too large"); + let mut bytes = (payload.len() as u32).to_be_bytes().to_vec(); + bytes.extend(payload); + Ok(bytes) +} + +fn read_frame(input: &mut impl Read) -> Result> { + let mut size = [0; 4]; + if input.read(&mut size[..1])? == 0 { + return Ok(None); + } + input.read_exact(&mut size[1..])?; + let size = u32::from_be_bytes(size) as usize; + anyhow::ensure!(size <= FRAME_LIMIT, "worker frame too large"); + let mut bytes = vec![0; size]; + input.read_exact(&mut bytes)?; + Ok(Some(serde_json::from_slice(&bytes)?)) +} + +pub(super) fn validate_backends(mounts: &[CatalogLibrary]) -> Result<()> { + anyhow::ensure!(!mounts.is_empty(), "catalog worker needs mounts"); + let mut backend = None; + for library in mounts.iter().filter(|m| m.uri.starts_with("s3://")) { + anyhow::ensure!( + library.access_key.is_some() && library.secret_key.is_some(), + "explicit S3 credentials required" + ); + let identity = ( + &library.endpoint, + &library.region, + &library.access_key, + &library.secret_key, + ); + if let Some(previous) = backend { + anyhow::ensure!( + previous == identity, + "datasets use different S3 credentials or endpoints; select a dataset explicitly" + ); + } + backend = Some(identity); + } + Ok(()) +} + +/// Called before main constructs any runtime or threads. Credentials arrive +/// only over stdin, and remain fixed for the lifetime of this process. +pub(crate) fn run() -> Result<()> { + let mut input = std::io::stdin().lock(); + let mut output = std::io::stdout().lock(); + let bootstrap: Bootstrap = read_frame(&mut input)?.context("missing worker bootstrap")?; + validate_backends(&bootstrap.mounts)?; + if let Some(library) = bootstrap.mounts.iter().find(|m| m.uri.starts_with("s3://")) { + apply_library_env(library); + } + let mounts = bootstrap + .mounts + .iter() + .map(|m| persisting_pchronicle::storage::DatasetMount::new(&m.name, &m.uri)) + .collect::>>()?; + let runtime = tokio::runtime::Builder::new_multi_thread() + .worker_threads(2) + .max_blocking_threads(8) + .enable_all() + .build()?; + let warehouse = runtime.block_on(super::PreparedWarehouse::prepare_query_worker( + super::ChronicleServerConfig::mounted(mounts)?, + ))?; + output.write_all(&encode(&true)?)?; + output.flush()?; + while let Some(job) = read_frame::(&mut input)? { + anyhow::ensure!( + job.body.len() <= 1024 * 1024, + "worker request body too large" + ); + let result = runtime.block_on(async { + let mut builder = axum::http::Request::builder() + .method(job.method.as_str()) + .uri(job.uri); + for (name, value) in job.headers { + builder = builder.header(name, value); + } + let response = warehouse + .router() + .oneshot(builder.body(axum::body::Body::from(job.body))?) + .await?; + let status = response.status().as_u16(); + let headers = response + .headers() + .iter() + .filter(|(name, _)| { + matches!( + name.as_str(), + "content-type" + | "server-timing" + | "x-request-id" + | "cache-control" + | "retry-after" + | "content-disposition" + ) + }) + .filter_map(|(name, value)| { + value + .to_str() + .ok() + .map(|v| (name.to_string(), v.to_owned())) + }) + .collect(); + let body = axum::body::to_bytes(response.into_body(), 8 * 1024 * 1024) + .await? + .to_vec(); + Ok::<_, anyhow::Error>(WorkerResponse { + status, + headers, + body, + }) + })?; + output.write_all(&encode(&result)?)?; + output.flush()?; + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn frames_reject_truncation_and_oversize_and_preserve_boundaries() { + let mut bytes = encode(&true).unwrap(); + bytes.extend(encode(&false).unwrap()); + let mut input = bytes.as_slice(); + assert_eq!(read_frame::(&mut input).unwrap(), Some(true)); + assert_eq!(read_frame::(&mut input).unwrap(), Some(false)); + assert_eq!(read_frame::(&mut input).unwrap(), None); + assert!(read_frame::(&mut &bytes[..3]).is_err()); + assert!(read_frame::(&mut &((FRAME_LIMIT + 1) as u32).to_be_bytes()[..]).is_err()); + } + + #[tokio::test] + async fn pool_bounds_admission_and_never_evicts_queued_scopes() { + let pool = WorkerPool::default(); + let permits: Vec<_> = (0..MAX_REQUESTS).map(|_| pool.admit().unwrap()).collect(); + assert!(pool.admit().is_err()); + drop(permits); + assert!(pool.admit().is_ok()); + let mut slots = Vec::new(); + for index in 0..MAX_WORKERS { + slots.push(pool.slot(&index.to_string()).await.unwrap()); + } + assert!(Arc::ptr_eq(&slots[0], &pool.slot("0").await.unwrap())); + assert!(pool.slot("overflow").await.is_err()); + slots.remove(0); + assert!(pool.slot("replacement").await.is_ok()); + assert_eq!(pool.slots.lock().await.len(), MAX_WORKERS); + } + + #[cfg(unix)] + #[tokio::test] + async fn cancelling_ipc_kills_worker_and_clears_slot() { + let home = tempfile::tempdir().unwrap(); + let mut child = tokio::process::Command::new("/bin/sh") + .args(["-c", "exec sleep 30"]) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .kill_on_drop(true) + .spawn() + .unwrap(); + let pid = child.id().unwrap() as i32; + let input = child.stdin.take().unwrap(); + let output = child.stdout.take().unwrap(); + let pool = WorkerPool::default(); + let slot = pool.slot("test").await.unwrap(); + *slot.lock().await = Some(Worker { + child, + input, + output, + _home: home, + }); + let result = tokio::time::timeout( + Duration::from_millis(30), + pool.execute( + "test".into(), + Vec::new(), + WorkerRequest { + method: "GET".into(), + uri: "/api/health".into(), + headers: Vec::new(), + body: Vec::new(), + }, + ), + ) + .await; + assert!(result.is_err()); + assert!(slot.lock().await.is_none()); + tokio::time::timeout(Duration::from_secs(5), async { + loop { + // Signal zero observes process existence without sending a signal. + if unsafe { libc::kill(pid, 0) } == -1 { + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .unwrap(); + } + + #[test] + fn exec_environment_excludes_ambient_credentials() { + let home = tempfile::tempdir().unwrap(); + let cmd = command(PathBuf::from("pchronicle"), home.path(), home.path()); + let env: HashMap<_, _> = cmd.as_std().get_envs().collect(); + for key in [ + "AWS_ACCESS_KEY_ID", + "AWS_SECRET_ACCESS_KEY", + "AWS_SESSION_TOKEN", + "AWS_PROFILE", + "AWS_WEB_IDENTITY_TOKEN_FILE", + "AWS_CONTAINER_CREDENTIALS_FULL_URI", + ] { + assert!(!env.contains_key(std::ffi::OsStr::new(key))); + } + assert_eq!( + env[std::ffi::OsStr::new("HOME")], + Some(home.path().as_os_str()) + ); + assert_eq!( + env[std::ffi::OsStr::new("AWS_EC2_METADATA_DISABLED")], + Some(std::ffi::OsStr::new("true")) + ); + } + + #[test] + fn mixed_s3_credentials_are_rejected_instead_of_using_first_key() { + let library = CatalogLibrary { + name: "one".into(), + uri: "s3://bucket/one".into(), + endpoint: None, + region: None, + access_key: Some("ak".into()), + secret_key: Some("sk".into()), + }; + let mut other = library.clone(); + other.name = "two".into(); + other.uri = "s3://bucket/two".into(); + assert!(validate_backends(&[library.clone(), other.clone()]).is_ok()); + other.secret_key = Some("another-secret".into()); + let error = validate_backends(&[library, other]) + .unwrap_err() + .to_string(); + assert!(!error.contains("another-secret")); + assert!(error.contains("select a dataset")); + } +} diff --git a/crates/persisting-pchronicle-cli/src/server/explorer.rs b/crates/persisting-pchronicle-cli/src/server/explorer.rs index 1fc80c7cf..d142fdf81 100644 --- a/crates/persisting-pchronicle-cli/src/server/explorer.rs +++ b/crates/persisting-pchronicle-cli/src/server/explorer.rs @@ -37,6 +37,10 @@ pub(crate) struct CatalogTree { pub(crate) prefix: String, pub(crate) run_count: usize, pub(crate) failed_count: usize, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub(crate) dataset_count: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub(crate) trajectory_count: Option, #[serde(skip_serializing_if = "Option::is_none")] pub(crate) ready_sources: Option, #[serde(skip_serializing_if = "Option::is_none")] @@ -56,6 +60,10 @@ pub(crate) struct CatalogTreeChild { pub(crate) path: String, pub(crate) run_count: usize, pub(crate) failed_count: usize, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub(crate) dataset_count: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub(crate) trajectory_count: Option, #[serde(skip_serializing_if = "Option::is_none")] pub(crate) total_tokens: Option, #[serde(default, skip_serializing_if = "Vec::is_empty")] @@ -90,6 +98,8 @@ pub(crate) fn catalog_tree_from_mounts( path: dataset.mount.name.clone(), run_count, failed_count, + dataset_count: Some(1), + trajectory_count: Some(run_count), total_tokens: None, entries: Vec::new(), } @@ -151,6 +161,8 @@ pub(crate) fn catalog_tree_from_path_list( .failed_count .and_then(|count| usize::try_from(count).ok()) .unwrap_or(0), + dataset_count: None, + trajectory_count: None, total_tokens: None, entries: Vec::new(), } diff --git a/crates/persisting-pchronicle-cli/src/server/mod.rs b/crates/persisting-pchronicle-cli/src/server/mod.rs index 3ce0dc794..1c3676663 100644 --- a/crates/persisting-pchronicle-cli/src/server/mod.rs +++ b/crates/persisting-pchronicle-cli/src/server/mod.rs @@ -3,6 +3,7 @@ mod acceleration; mod asset; pub(crate) mod catalog; +pub(crate) mod catalog_worker; mod explorer; mod physical; pub(crate) mod problem; @@ -64,6 +65,8 @@ struct AppState { live_reads: bool, catalog_acl: Option>, catalog_query_worker: bool, + catalog_workers: Arc, + browse_mounts: Arc>, browse: Arc>, scoped_queries: Arc, } @@ -207,6 +210,7 @@ fn app_state_with_catalog_refresh_interval( config: ChronicleServerConfig, catalog_refresh_interval: Duration, ) -> AppState { + let browse_mounts = config.datasets.clone(); AppState { config: Arc::new(config), catalog: Arc::new(tokio::sync::RwLock::new(None)), @@ -216,6 +220,8 @@ fn app_state_with_catalog_refresh_interval( live_reads: false, catalog_acl: None, catalog_query_worker: false, + catalog_workers: Arc::new(catalog_worker::WorkerPool::default()), + browse_mounts: Arc::new(browse_mounts), browse: Arc::new(tokio::sync::OnceCell::new()), scoped_queries: Arc::new(query_admission::ScopedQueries::default()), } @@ -245,38 +251,32 @@ impl PreparedWarehouse { Ok(warehouse) } - /// Directory front-only mode: parent authenticates and dispatches query - /// workers. Retained for isolation tests; `serve --catalog-config` uses - /// [`Self::prepare_catalog`] (inline mounts) instead. - #[cfg_attr(not(test), allow(dead_code))] + #[cfg(test)] pub(crate) async fn prepare_catalog_front(acl: catalog::CatalogAcl) -> anyhow::Result { - let mut state = app_state(ChronicleServerConfig::front_only()); - state.catalog_acl = Some(Arc::new(acl)); - Ok(Self { state }) + Self::prepare_catalog(acl, ChronicleServerConfig::front_only()).await } - /// Mount every library from `catalog.toml` into the Warehouse process. - /// Directory ticket routes remain available when users exist; the data - /// plane serves in-process mounts instead of spawning query workers. - /// - /// Discovery runs in the background so `serve --listen` can accept - /// connections before large object prefixes finish classifying. + /// The listener authenticates and brokers credentials; only workers mount + /// datasets. Keep public UI configuration, never storage state, in the parent. pub(crate) async fn prepare_catalog( acl: catalog::CatalogAcl, - config: ChronicleServerConfig, + mut config: ChronicleServerConfig, ) -> anyhow::Result { - acl.apply_backend_env(); - anyhow::ensure!( - !config.datasets.is_empty(), - "catalog config needs at least one dataset" - ); + 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)); - let warehouse = Self { state }; - browse_coordinator(&warehouse.state).await; - // The UI is served by BrowseCoordinator. Build the accurate query - // snapshot lazily when an endpoint actually needs it. - Ok(warehouse) + // 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; + Ok(Self { state }) } pub(crate) async fn prepare_query_worker( @@ -285,7 +285,7 @@ impl PreparedWarehouse { let mut state = app_state(config); state.catalog_query_worker = true; let warehouse = Self { state }; - warehouse.install_initial_runtime().await?; + browse_coordinator(&warehouse.state).await; Ok(warehouse) } @@ -344,7 +344,7 @@ impl PreparedWarehouse { async fn browse_coordinator(state: &AppState) -> &ui_cache::BrowseCoordinator { state .browse - .get_or_init(|| ui_cache::BrowseCoordinator::start(state.config.datasets.clone())) + .get_or_init(|| ui_cache::BrowseCoordinator::start((*state.browse_mounts).clone())) .await } @@ -570,7 +570,7 @@ async fn current_catalog_for_runs( return Ok(runtime); } // A zero interval requests synchronous freshness (used by embedded callers). - if state.catalog_refresh_interval.is_zero() || state.live_reads || state.catalog_query_worker { + if state.catalog_refresh_interval.is_zero() || state.live_reads { let _refresh = state.catalog_refresh.lock().await; return Ok(rebuild_catalog_for_runs(state, runtime).await); } @@ -700,7 +700,7 @@ async fn load_run_summaries( }; let config = state.config.clone(); let query_scope = scope.clone(); - let cached_files = if file.is_none() && !state.live_reads && !state.catalog_query_worker { + let cached_files = if file.is_none() && !state.live_reads { browse_coordinator(state) .await .cached_source_paths( @@ -741,7 +741,7 @@ async fn load_run_summaries( summaries }; let started = Instant::now(); - let result = if state.live_reads || state.catalog_query_worker { + let result = if state.live_reads { state .scoped_queries .run(scope, || execute(false)) @@ -1375,21 +1375,34 @@ async fn explorer_tree( .map(str::trim) .filter(|value| !value.is_empty()); let prefix = query.prefix.as_deref().unwrap_or(""); + tracing::info!( + target: "pchronicle.serve", + dataset = ?dataset, + prefix, + catalog_worker = state.catalog_query_worker, + "explorer tree request" + ); let Some(name) = dataset else { let started = Instant::now(); let view = browse_coordinator(&state) .await - .roots(&state.config.datasets) + .roots(&state.browse_mounts) .await; metrics.record("browse", started); return Ok(Json(serde_json::to_value(view).unwrap())); }; - let Some(mount) = state - .config - .datasets - .iter() - .find(|mount| mount.name == name) - else { + 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() + .into_iter() + .find(|library| library.name == name) + }) { + owned_mount = DatasetMount::new(&library.name, &library.uri) + .map_err(|error| fail(&request_id, "explorer_tree", error))?; + &owned_mount + } else { return Ok(Json( serde_json::to_value(explorer::catalog_tree_from_path_list(name, prefix, &[])).unwrap(), )); diff --git a/crates/persisting-pchronicle-cli/src/server/request_log.rs b/crates/persisting-pchronicle-cli/src/server/request_log.rs index 13011a777..d077696a9 100644 --- a/crates/persisting-pchronicle-cli/src/server/request_log.rs +++ b/crates/persisting-pchronicle-cli/src/server/request_log.rs @@ -192,8 +192,15 @@ async fn attach_request_id( if let Ok(value) = HeaderValue::from_str(request_id) { parts.headers.insert("x-request-id", value); } - if let Ok(value) = HeaderValue::from_str(&metrics.server_timing(total_ms)) { - parts.headers.insert("server-timing", value); + let timing = if parts.headers.contains_key("server-timing") { + // Preserve stages returned by an isolated worker; parent wall time + // includes admission and IPC and must have a distinct metric name. + format!("parent_total;dur={total_ms}") + } else { + metrics.server_timing(total_ms) + }; + if let Ok(value) = HeaderValue::from_str(&timing) { + parts.headers.append("server-timing", value); } let is_json = parts .headers diff --git a/crates/persisting-pchronicle-cli/src/server/ui_cache.rs b/crates/persisting-pchronicle-cli/src/server/ui_cache.rs index 010535baf..1f50ed030 100644 --- a/crates/persisting-pchronicle-cli/src/server/ui_cache.rs +++ b/crates/persisting-pchronicle-cli/src/server/ui_cache.rs @@ -110,6 +110,7 @@ type Pending = Arc>>>; /// aborts it, including an in-flight list. No permanent process singleton. pub(crate) struct BrowseCoordinator { index: Arc, + manifests: Arc, pending: Pending, refresh: Arc>>, sender: mpsc::Sender, @@ -144,6 +145,12 @@ impl BrowseCoordinator { ); let manifests = Arc::new(ManifestCache::open(root.join(format!("manifest-{namespace}.lance"))).await); + tracing::info!( + target: "pchronicle.serve", + mounts = mounts.len(), + cache_root = %root.display(), + "catalog cache worker started" + ); let pending = Arc::new(Mutex::new(HashMap::new())); let refresh = Arc::new(Mutex::new(HashMap::new())); let (sender, receiver) = mpsc::channel(QUEUE_CAPACITY); @@ -157,6 +164,7 @@ impl BrowseCoordinator { )); Self { index, + manifests, pending, refresh, sender, @@ -218,11 +226,27 @@ impl BrowseCoordinator { } else { complete = false; } + let summary = self + .manifests + .summary(&format!("{}\0{}", mount.name, mount_fingerprint(mount))) + .await; + if let Some(child) = tree.children.iter_mut().find(|c| c.name == mount.name) { + child.run_count = summary.trajectories as usize; + child.dataset_count = Some(summary.datasets as usize); + child.trajectory_count = Some(summary.trajectories as usize); + } } tree.run_count = tree .children .iter() .fold(0usize, |sum, c| sum.saturating_add(c.run_count)); + tree.dataset_count = Some( + tree.children + .iter() + .map(|child| child.dataset_count.unwrap_or_default()) + .sum(), + ); + tree.trajectory_count = Some(tree.run_count); tree.failed_count = tree .children .iter() @@ -244,10 +268,16 @@ impl BrowseCoordinator { pub(crate) async fn tree(&self, mount: &DatasetMount, prefix: &str) -> Result { let key = TreeKey::new(mount, prefix)?; + tracing::info!( + target: "pchronicle.serve", + dataset = %key.dataset, + prefix = %key.prefix, + "browse tree request" + ); let existing = self.index.values.read().await.get(&key).cloned(); if let Some(entry) = existing { let _ = self.enqueue(&key, None); - return Ok(self.snapshot(&key, entry)); + return Ok(self.snapshot_with_summary(&key, entry).await); } let (reply, wait) = oneshot::channel(); self.enqueue(&key, Some(reply))?; @@ -265,7 +295,7 @@ impl BrowseCoordinator { .get(&key) .cloned() .context("browse refresh produced no view")?; - Ok(self.snapshot(&key, entry)) + Ok(self.snapshot_with_summary(&key, entry).await) } fn enqueue(&self, key: &TreeKey, reply: Option) -> Result<()> { @@ -315,6 +345,132 @@ impl BrowseCoordinator { }, } } + + async fn snapshot_with_summary(&self, key: &TreeKey, mut entry: IndexEntry) -> BrowseSnapshot { + let current_manifest_key = manifest_key(key, None); + let summary = self.manifests.summary_under(¤t_manifest_key).await; + entry.tree.dataset_count = Some(summary.datasets as usize); + entry.tree.trajectory_count = Some(summary.trajectories as usize); + entry.tree.run_count = summary.trajectories as usize; + let cached_trees = self + .index + .values + .read() + .await + .values() + .map(|entry| entry.tree.clone()) + .collect::>(); + let (cached_datasets, cached_trajectories) = cached_leaf_summary(&cached_trees, ""); + if cached_datasets > 0 { + entry.tree.dataset_count = Some(cached_datasets); + entry.tree.trajectory_count = Some(cached_trajectories); + entry.tree.run_count = cached_trajectories; + } + for child in &mut entry.tree.children { + if child.kind != "dir" { + continue; + } + let child_prefix = if key.prefix.is_empty() + || child.path == key.prefix + || child.path.starts_with(&format!("{}/", key.prefix)) + { + child.path.clone() + } else { + format!("{}/{}", key.prefix, child.path) + }; + let summary = self + .manifests + .summary_under(&manifest_key(key, Some(&child_prefix))) + .await; + child.dataset_count = (summary.datasets > 0).then_some(summary.datasets as usize); + child.trajectory_count = + (summary.trajectories > 0).then_some(summary.trajectories as usize); + let child_key = TreeKey { + dataset: key.dataset.clone(), + uri_fingerprint: key.uri_fingerprint.clone(), + prefix: child_prefix.trim_matches('/').to_owned(), + }; + // A directory can be discovered after the background walk started. + // Schedule it here so its descendant manifests become available to + // the next render without making the request wait on storage I/O. + let _ = self.enqueue(&child_key, None); + let cached = self.index.values.read().await.get(&child_key).cloned(); + tracing::info!( + target: "pchronicle.serve", + dataset = %key.dataset, + parent_prefix = %key.prefix, + child_prefix = %child.path, + manifest_datasets = summary.datasets, + manifest_trajectories = summary.trajectories, + projection_hit = cached.is_some(), + "catalog directory summary" + ); + if let Some(cached) = cached { + if let Some(count) = cached.tree.dataset_count.filter(|count| *count > 0) { + child.dataset_count = Some(count); + } + if let Some(count) = cached + .tree + .trajectory_count + .or(Some(cached.tree.run_count)) + .filter(|count| *count > 0) + { + child.trajectory_count = Some(count); + } + } + let (datasets, trajectories) = cached_leaf_summary(&cached_trees, &child_prefix); + if datasets > 0 { + child.dataset_count = Some(datasets); + child.trajectory_count = Some(trajectories); + } + } + let child_summary = + entry + .tree + .children + .iter() + .fold((0usize, 0usize), |(datasets, trajectories), child| { + let datasets = datasets.saturating_add( + child + .dataset_count + .unwrap_or_else(|| (child.kind == "dataset") as usize), + ); + let trajectories = trajectories + .saturating_add(child.trajectory_count.unwrap_or(child.run_count)); + (datasets, trajectories) + }); + entry.tree.dataset_count = Some((summary.datasets as usize).max(child_summary.0)); + entry.tree.trajectory_count = Some((summary.trajectories as usize).max(child_summary.1)); + entry.tree.run_count = entry.tree.trajectory_count.unwrap_or_default(); + self.snapshot(key, entry) + } +} + +fn manifest_key(key: &TreeKey, child_prefix: Option<&str>) -> String { + let prefix = child_prefix.unwrap_or(&key.prefix); + if prefix.is_empty() { + format!("{}\0{}", key.dataset, key.uri_fingerprint) + } else { + format!("{}\0{}\0{}", key.dataset, key.uri_fingerprint, prefix) + } +} + +fn cached_leaf_summary(trees: &[CatalogTree], prefix: &str) -> (usize, usize) { + let mut leaves = HashMap::::new(); + for tree in trees { + for child in &tree.children { + if child.kind == "file" + && child.data_type != "other" + && (child.path == prefix || child.path.starts_with(&format!("{prefix}/"))) + { + leaves.insert(child.path.clone(), child.run_count); + } + } + } + ( + leaves.len(), + leaves.values().copied().fold(0usize, usize::saturating_add), + ) } async fn run_worker( @@ -336,8 +492,10 @@ async fn run_worker( loop { let key = tokio::select! { biased; + _ = std::future::ready(()), if !background.is_empty() => background.pop_front().unwrap(), request = receiver.recv() => match request { Some(key) => key, None => break }, _ = interval.tick() => { + tracing::info!(target: "pchronicle.serve", mounts = mounts.len(), "browse manifest refresh round started"); if background.is_empty() { visited.clear(); background.extend(mounts.iter().filter_map(|m| TreeKey::new(m, "").ok())); @@ -345,7 +503,6 @@ async fn run_worker( } continue; } - _ = std::future::ready(()), if !background.is_empty() => background.pop_front().unwrap(), }; let Some(mount) = mounts .iter() @@ -377,7 +534,11 @@ async fn run_worker( .await .context("browse path unavailable")?; } - let manifest_key = format!("{}\0{}\0{}", key.dataset, key.uri_fingerprint, key.prefix); + let manifest_key = if key.prefix.is_empty() { + format!("{}\0{}", key.dataset, key.uri_fingerprint) + } else { + format!("{}\0{}\0{}", key.dataset, key.uri_fingerprint, key.prefix) + }; let listing = tokio::time::timeout( Duration::from_secs(20), manifests.refresh(manifest_key, &location, &key.prefix), @@ -403,6 +564,7 @@ async fn run_worker( .await; let outcome = match result { Ok(()) => { + tracing::info!(target: "pchronicle.serve", dataset = %key.dataset, prefix = %key.prefix, "browse manifest level refreshed"); states.lock().unwrap().insert( key.clone(), RefreshState { @@ -546,6 +708,70 @@ mod tests { DatasetMount::new("test", path.to_string_lossy()).unwrap() } + #[tokio::test] + async fn directory_statistics_use_cached_descendants_without_child_projections() { + let temp = tempfile::tempdir().unwrap(); + let source = temp.path().join("source"); + for (path, count) in [ + ("codex2", 0), + ("codex3", 0), + ("nested/nested1/codex2", 354), + ("nested/nested2/codex2", 354), + ] { + let leaf = source.join(path); + std::fs::create_dir_all(&leaf).unwrap(); + persisting_pchronicle::storage::write_compact_jsonl_manifest(&leaf, 1, count).unwrap(); + } + let mount = mount(&source); + // Keep the scheduler idle: only manifest observations supply the response. + let browse = BrowseCoordinator::start_at(Vec::new(), temp.path().join("cache")).await; + browse.task.abort(); + let root_key = TreeKey::new(&mount, "").unwrap(); + browse + .manifests + .refresh_mount( + &manifest_key(&root_key, None), + &DatasetLocation::parse(&mount.uri).unwrap(), + ) + .await + .unwrap(); + std::fs::remove_dir_all(&source).unwrap(); + assert!(browse.index.values.read().await.is_empty()); + for (prefix, datasets) in [("", 4), ("nested", 2)] { + let key = TreeKey::new(&mount, prefix).unwrap(); + let listing = browse + .manifests + .get(&manifest_key(&key, None)) + .await + .unwrap(); + let view = browse + .snapshot_with_summary( + &key, + IndexEntry { + tree: catalog_tree_from_path_list(&mount.name, prefix, &listing.entries), + generation: String::new(), + observed_at: now(), + }, + ) + .await; + let json = serde_json::to_value(view).unwrap(); + assert_eq!(json["dataset_count"], datasets); + assert_eq!(json["trajectory_count"], 708); + for child in json["children"].as_array().unwrap() { + if child["kind"] == "dir" { + assert_eq!( + child["dataset_count"], + if prefix.is_empty() { 2 } else { 1 } + ); + assert_eq!( + child["trajectory_count"], + if prefix.is_empty() { 708 } else { 354 } + ); + } + } + } + } + #[tokio::test] async fn corrupt_cache_is_rebuilt_and_persisted() { let temp = tempfile::tempdir().unwrap(); diff --git a/crates/persisting-pchronicle-cli/src/tests.rs b/crates/persisting-pchronicle-cli/src/tests.rs index 537bc0b1d..881da7157 100644 --- a/crates/persisting-pchronicle-cli/src/tests.rs +++ b/crates/persisting-pchronicle-cli/src/tests.rs @@ -391,7 +391,7 @@ fn command_tree_contains_the_product_commands() { serve_help.write_long_help(&mut help).unwrap(); let help = String::from_utf8(help).unwrap(); assert!(help.contains("pchronicle serve catalog"), "{help}"); - assert!(help.contains("Mounts every [datasets.*] entry"), "{help}"); + assert!(help.contains("Authenticate API requests"), "{help}"); } #[test] diff --git a/crates/persisting-pchronicle-cli/tests/catalog_worker_contract.rs b/crates/persisting-pchronicle-cli/tests/catalog_worker_contract.rs new file mode 100644 index 000000000..add3c50c5 --- /dev/null +++ b/crates/persisting-pchronicle-cli/tests/catalog_worker_contract.rs @@ -0,0 +1,228 @@ +//! Exercise the actual exec boundary, not a router with a test-only worker. +use anyhow::{Context, Result}; +use serde_json::{Value, json}; +use std::process::Stdio; +use std::time::Duration; +use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader}; + +#[tokio::test] +async fn catalog_server_authenticates_each_request_and_isolates_grants() -> Result<()> { + let home = tempfile::tempdir()?; + for name in ["left", "right"] { + std::fs::create_dir(home.path().join(name))?; + let mut trajectory: Value = serde_json::from_str(include_str!( + "../../../examples/data/atif/support-ticket.json" + ))?; + trajectory["session_id"] = json!(format!("{name}-session")); + std::fs::write( + home.path().join(name).join("trajectory.atif.json"), + serde_json::to_vec(&trajectory)?, + )?; + } + let config = home.path().join("catalog.toml"); + std::fs::write( + &config, + format!( + r#" +[datasets.left] +uri = "{}" +[datasets.right] +uri = "{}" +[users.alice] +access_key = "alice-ak" +secret_key = "alice-sk" +[users.bob] +access_key = "bob-ak" +secret_key = "bob-sk" +[[grants]] +user = "alice" +dataset = "left" +permissions = ["read"] +[[grants]] +user = "bob" +dataset = "right" +permissions = ["read"] +"#, + home.path().join("left").display(), + home.path().join("right").display() + ), + )?; + let cache = home.path().join("cache"); + let server_log = home.path().join("server.stderr"); + let mut server = tokio::process::Command::new(env!("CARGO_BIN_EXE_pchronicle")) + .args(["serve", "--catalog-config"]) + .arg(&config) + .args(["--listen", "127.0.0.1:0"]) + .env("PCHRONICLE_CACHE_DIR", &cache) + .env("AWS_ACCESS_KEY_ID", "ambient-must-not-be-used") + .env("AWS_SECRET_ACCESS_KEY", "ambient-secret") + .stdout(Stdio::piped()) + .stderr(std::fs::File::create(&server_log)?) + .kill_on_drop(true) + .spawn()?; + let mut stdout = BufReader::new(server.stdout.take().context("server stdout")?); + let mut line = String::new(); + tokio::time::timeout(Duration::from_secs(60), stdout.read_line(&mut line)) + .await + .context("wait for catalog listener readiness")??; + let ready: Value = serde_json::from_str(&line).with_context(|| { + format!( + "server readiness; stderr: {}", + std::fs::read_to_string(&server_log).unwrap_or_default() + ) + })?; + let endpoint = ready["warehouse_endpoint"] + .as_str() + .context("warehouse endpoint")?; + let client = reqwest::Client::builder() + .timeout(Duration::from_secs(70)) + .build()?; + let url = format!("http://{endpoint}/api/explorer/tree"); + assert_anonymous_access_is_scoped(&client, endpoint).await?; + // Same keep-alive client alternates identities; identity belongs to the + // request, never to the TCP connection or the previous worker. + for (user, dataset) in [("alice", "left"), ("bob", "right"), ("alice", "left")] { + let response = client + .get(&url) + .header("x-pchronicle-access-key", format!("{user}-ak")) + .header("x-pchronicle-secret-key", format!("{user}-sk")) + .send() + .await?; + assert_eq!(response.status(), 200); + let timings = response + .headers() + .get_all("server-timing") + .iter() + .map(|v| v.to_str().unwrap_or("")) + .collect::>() + .join(","); + assert!(timings.contains("parent_total"), "{timings}"); + assert!(timings.contains("total;"), "{timings}"); + let body: Value = response.json().await?; + assert_eq!(body["children"].as_array().map(Vec::len), Some(1)); + assert_eq!(body["children"][0]["name"], dataset); + let query = client + .post(format!("http://{endpoint}/api/query/evidence")) + .header("x-pchronicle-access-key", format!("{user}-ak")) + .header("x-pchronicle-secret-key", format!("{user}-sk")) + .json(&json!({"sql":"SELECT session_id FROM runs", "max_rows":10, "max_bytes":4096})) + .send() + .await?; + assert_eq!(query.status(), 200); + let query: Value = query.json().await?; + assert_eq!( + query["rows"], + json!([{"session_id":format!("{dataset}-session")}]), + "{query}" + ); + } + let response = client + .get(format!("{url}?dataset=right")) + .header("x-pchronicle-access-key", "alice-ak") + .header("x-pchronicle-secret-key", "alice-sk") + .send() + .await?; + assert_eq!(response.status(), 404); + assert_anonymous_access_is_scoped(&client, endpoint).await?; + assert_eq!(std::fs::read_dir(cache.join("workers"))?.count(), 2); + server.kill().await?; + Ok(()) +} + +async fn assert_anonymous_access_is_scoped(client: &reqwest::Client, endpoint: &str) -> Result<()> { + let tree = format!("http://{endpoint}/api/explorer/tree"); + let response = client.get(&tree).send().await?; + assert_eq!(response.status(), 200); + let body: Value = response.json().await?; + // This catalog has no wildcard grants: private mounts must stay hidden, + // including after this keep-alive client has used an authenticated worker. + assert_eq!(body["children"], json!([]), "{body}"); + for dataset in ["left", "right"] { + assert_eq!( + client + .get(&tree) + .query(&[("dataset", dataset)]) + .send() + .await? + .status(), + 401 + ); + } + assert_eq!( + client + .get(&tree) + .header("x-pchronicle-access-key", "alice-ak") + .header("x-pchronicle-secret-key", "wrong-secret") + .send() + .await? + .status(), + 401 + ); + assert_eq!( + client + .post(format!("http://{endpoint}/api/query/evidence")) + .json(&json!({"sql":"SELECT session_id FROM runs", "max_rows":10, "max_bytes":4096})) + .send() + .await? + .status(), + 401 + ); + Ok(()) +} + +async fn frame(input: &mut tokio::process::ChildStdin, value: Value) -> Result<()> { + let bytes = serde_json::to_vec(&value)?; + input.write_u32(bytes.len() as u32).await?; + input.write_all(&bytes).await?; + input.flush().await?; + Ok(()) +} + +async fn read_frame(output: &mut tokio::process::ChildStdout) -> Result { + let size = tokio::time::timeout(Duration::from_secs(30), output.read_u32()).await?? as usize; + anyhow::ensure!(size < 1024 * 1024, "unexpected test response size"); + let mut bytes = vec![0; size]; + output.read_exact(&mut bytes).await?; + Ok(serde_json::from_slice(&bytes)?) +} + +#[tokio::test] +async fn exec_worker_handles_multiple_frames_and_exits_on_eof() -> Result<()> { + let home = tempfile::tempdir()?; + let mut child = tokio::process::Command::new(env!("CARGO_BIN_EXE_pchronicle")) + .args(["serve", "--catalog-query-worker"]) + .env_clear() + .env("HOME", home.path()) + .env("PCHRONICLE_CACHE_DIR", home.path().join("cache")) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::null()) + .kill_on_drop(true) + .spawn()?; + let pid = child.id(); + let mut input = child.stdin.take().context("worker stdin")?; + let mut output = child.stdout.take().context("worker stdout")?; + frame( + &mut input, + json!({"mounts":[{"name":"local", "uri":home.path().join("dataset")}]}), + ) + .await?; + assert_eq!(read_frame(&mut output).await?, true); + for _ in 0..2 { + frame( + &mut input, + json!({"method":"GET", "uri":"/api/health", "headers":[], "body":[]}), + ) + .await?; + assert_eq!(read_frame(&mut output).await?["status"], 200); + assert_eq!(child.id(), pid); + assert!(child.try_wait()?.is_none()); + } + drop(input); + assert!( + tokio::time::timeout(Duration::from_secs(10), child.wait()) + .await?? + .success() + ); + Ok(()) +} diff --git a/crates/persisting-pchronicle/src/store/catalog/manifest_cache.rs b/crates/persisting-pchronicle/src/store/catalog/manifest_cache.rs index 6ecf404b5..0c1734365 100644 --- a/crates/persisting-pchronicle/src/store/catalog/manifest_cache.rs +++ b/crates/persisting-pchronicle/src/store/catalog/manifest_cache.rs @@ -3,7 +3,7 @@ //! This module deliberately knows nothing about DataFusion. It is the single //! boundary used by navigational callers that can tolerate a stale snapshot. -use std::collections::VecDeque; +use std::collections::{HashSet, VecDeque}; use std::path::PathBuf; use std::sync::Arc; use std::time::Duration; @@ -104,6 +104,7 @@ impl ManifestCache { self.disk .upsert(&key, &serde_json::to_value(&listing)?, &[]) .await?; + tracing::info!(target: "pchronicle.serve", %key, prefix, entries = listing.entries.len(), "manifest cache updated"); Ok(listing) } @@ -145,13 +146,30 @@ impl ManifestCache { /// Aggregate the currently cached manifest observations for one mount. /// This never performs I/O and is therefore safe for UI rendering. pub async fn summary(&self, key_prefix: &str) -> LocationSummary { + self.summary_under(key_prefix).await + } + + /// Aggregate every cached manifest at `key_prefix` and below it. + pub async fn summary_under(&self, key_prefix: &str) -> LocationSummary { let values = self.values.read().await; + let mut seen = HashSet::new(); values .iter() - .filter(|(key, _)| key == &key_prefix || key.starts_with(&format!("{key_prefix}\0"))) + .filter(|(key, _)| { + // NUL separates the mount identity from the relative path; + // slashes separate directories inside that path. + key.strip_prefix(key_prefix).is_some_and(|suffix| { + suffix.is_empty() || suffix.starts_with('\0') || suffix.starts_with('/') + }) + }) .fold(LocationSummary::default(), |mut total, (_, listing)| { for entry in &listing.entries { - if matches!(entry.kind, crate::store::PathListKind::Dataset) { + let is_dataset = matches!(entry.kind, crate::store::PathListKind::Dataset) + || (matches!(entry.kind, crate::store::PathListKind::File) + && entry.format.as_deref().is_none_or(|format| { + matches!(format, "storyline-lance" | "compact-jsonl/v1") + })); + if is_dataset && seen.insert(entry.path.clone()) { total.datasets += 1; total.trajectories += entry.record_count.unwrap_or_default(); } @@ -164,7 +182,9 @@ impl ManifestCache { /// shallow paths are published before deeper paths. pub async fn refresh_mount(&self, key_prefix: &str, location: &DatasetLocation) -> Result<()> { let _guard = self.refresh_gate.lock().await; + tracing::info!(target: "pchronicle.serve", %key_prefix, "manifest cache refresh started"); let mut queue = VecDeque::from([String::new()]); + let mut refreshed = 0usize; while let Some(prefix) = queue.pop_front() { let listing = ManifestListing { entries: location.list(&prefix).await?, @@ -182,6 +202,7 @@ impl ManifestCache { self.disk .upsert(&key, &serde_json::to_value(&listing)?, &[]) .await?; + refreshed += 1; for child in listing .entries .iter() @@ -191,6 +212,7 @@ impl ManifestCache { } tokio::task::yield_now().await; } + tracing::info!(target: "pchronicle.serve", %key_prefix, refreshed, "manifest cache refresh finished"); Ok(()) } @@ -243,6 +265,90 @@ mod tests { trajectories: 3 } ); + cache.values.write().await.insert( + "mount\0nested".into(), + ManifestListing { + entries: vec![PathListEntry { + name: "b".into(), + path: "nested/b".into(), + kind: PathListKind::File, + format: Some("storyline-lance".into()), + record_count: Some(4), + failed_count: None, + }], + observed_at: 0, + }, + ); + assert_eq!(cache.summary_under("mount").await.trajectories, 7); + } + + #[tokio::test] + async fn summary_under_aggregates_slash_descendants_from_cache() { + let temp = tempfile::tempdir().unwrap(); + let source = temp.path().join("source"); + for (path, count) in [ + ("codex2", 0), + ("codex3", 0), + ("nested/nested1/codex2", 354), + ("nested/nested2/codex2", 354), + ] { + let leaf = source.join(path); + std::fs::create_dir_all(&leaf).unwrap(); + crate::store::catalog::manifest::write_compact_jsonl_manifest(&leaf, 1, count).unwrap(); + } + let location = DatasetLocation::parse(source.to_str().unwrap()).unwrap(); + let cache_path = temp.path().join("manifest.lance"); + let cache = ManifestCache::open(cache_path.clone()).await; + let mount = "rfs\0fingerprint"; + cache.refresh_mount(mount, &location).await.unwrap(); + assert!( + cache + .get("rfs\0fingerprint\0nested/nested1") + .await + .is_some() + ); + drop(cache); + std::fs::remove_dir_all(&source).unwrap(); + + // No source I/O or projection merging: aggregate persisted observations. + let cache = ManifestCache::open(cache_path).await; + for (prefix, datasets, trajectories) in [ + ("rfs\0fingerprint", 4, 708), + ("rfs\0fingerprint\0nested", 2, 708), + ("rfs\0fingerprint\0nested/nested1", 1, 354), + ("rfs\0fingerprint\0missing", 0, 0), + ] { + assert_eq!( + cache.summary_under(prefix).await, + LocationSummary { + datasets, + trajectories + }, + "{prefix:?}" + ); + } + let distractor = cache.get("rfs\0fingerprint\0nested/nested1").await.unwrap(); + let mut distractor = distractor; + distractor.entries[0].path = "unrelated/leaf".into(); + distractor.entries[0].record_count = Some(999); + for key in [ + "rfs\0fingerprint\0nested-other/child", + "rfs\0different\0nested/child", + "other\0fingerprint\0nested/child", + ] { + cache + .values + .write() + .await + .insert(key.into(), distractor.clone()); + } + assert_eq!( + cache.summary_under("rfs\0fingerprint\0nested").await, + LocationSummary { + datasets: 2, + trajectories: 708 + } + ); } #[tokio::test] diff --git a/crates/persisting-pvisor/tests/macos_safe.rs b/crates/persisting-pvisor/tests/macos_safe.rs index 93d69917b..378ec8720 100644 --- a/crates/persisting-pvisor/tests/macos_safe.rs +++ b/crates/persisting-pvisor/tests/macos_safe.rs @@ -73,6 +73,12 @@ fn safe_profile_stages_reviews_and_applies_on_macos() { .arg(&outside) .arg(&outside_secret); let output = command.output().expect("run macOS safe profile"); + if !output.status.success() + && String::from_utf8_lossy(&output.stderr).contains("file system is not available") + { + eprintln!("skipping macOS safe-profile smoke test: macFUSE is installed but unavailable"); + return; + } assert!( output.status.success(), "stdout:\n{}\nstderr:\n{}", @@ -217,6 +223,12 @@ raise SystemExit(0 if inet_code in denied and loopback_code == 0 and host_code i .arg(&outside_socket) .output() .expect("run macOS deny-all profile"); + if !output.status.success() + && String::from_utf8_lossy(&output.stderr).contains("file system is not available") + { + eprintln!("skipping macOS deny-all test: macFUSE is installed but unavailable"); + return; + } assert!( output.status.success(), "stdout:\n{}\nstderr:\n{}", diff --git a/docs/src/en/pchronicle/design/architecture.md b/docs/src/en/pchronicle/design/architecture.md index 9afedc87a..fd13fc655 100644 --- a/docs/src/en/pchronicle/design/architecture.md +++ b/docs/src/en/pchronicle/design/architecture.md @@ -149,10 +149,12 @@ Snapshot before switching readers. Dataset tables prune by Source before opening matching fixed versions; caches and routing indexes are tied to that Snapshot generation. -With `--catalog-config`, Warehouse mounts every `[datasets.*]` library from the -Directory ACL file (same data plane as positional mounts) and also serves -Directory list/ticket routes for `catalog://` pins. Backend S3 endpoint, -region, and keys from the file are applied before stores open. After a CLI +With `--catalog-config`, the listener authenticates each request and dispatches +it to a bounded exec worker pool. Each worker has immutable user/grant/backend +credential scope and private caches; the parent never mounts datasets. Backend +credentials arrive over private IPC before runtime threads start. Workers retain +the server's OS identity; process separation is not a filesystem sandbox. +Directory list/ticket routes remain available for `catalog://` pins. After a CLI ticket, the client opens the ticket `uri` (a path) with storage credentials. That is platform addressing over paths, not a new Dataset kind. diff --git a/docs/src/en/pchronicle/guides/serve.md b/docs/src/en/pchronicle/guides/serve.md index b3ab9ff15..7fa94be38 100644 --- a/docs/src/en/pchronicle/guides/serve.md +++ b/docs/src/en/pchronicle/guides/serve.md @@ -69,6 +69,7 @@ pchronicle serve --catalog-config catalog.toml --listen 127.0.0.1:8081 users. `serve catalog dataset add|remove|list` rewrites libraries without starting HTTP. `serve catalog issue` writes a user with empty grants and prints the secret once on stdout; `grant` / `revoke` change which library names that +Use `*` as NAME to grant or revoke a dataset for every current user. New users do not inherit past wildcard grants; rerun the command after creating them. user may open. Restart serve after editing the file. `pchronicle serve --catalog-config` mounts **every** library in the file into diff --git a/docs/src/en/pchronicle/guides/ui.md b/docs/src/en/pchronicle/guides/ui.md index 6fbfbf033..b2aeda31b 100644 --- a/docs/src/en/pchronicle/guides/ui.md +++ b/docs/src/en/pchronicle/guides/ui.md @@ -151,9 +151,9 @@ untrusted or shared browser profile. Clearing this site's browser data also clears the setting. Assistant is labeled **Read-only · selected run data** and does not rewrite the Dataset. -When `pchronicle serve --catalog-config` is used, every library in the ACL file -is already mounted for local browsing. Open **Keys** on the left rail if you -need Directory user access/secret headers for authenticated Directory flows. +With `pchronicle serve --catalog-config`, data APIs require authentication and +show only datasets granted to the user. Open **Keys** on the left rail and +enter the Directory user access key and secret key before browsing data. Those values are stored in `localStorage` and sent to this pchronicle server as `x-pchronicle-access-key` and `x-pchronicle-secret-key` on data requests. They authorize which Directory paths this browser may open; they are not the diff --git a/docs/src/en/pchronicle/reference/cases-platform.md b/docs/src/en/pchronicle/reference/cases-platform.md index 24545a4b6..752498ec3 100644 --- a/docs/src/en/pchronicle/reference/cases-platform.md +++ b/docs/src/en/pchronicle/reference/cases-platform.md @@ -98,6 +98,6 @@ Platform checks: - user, dataset, and grant edits are deterministic; - backend object-store keys stay in the catalog file / ticket path, not in `dataset list` output; -- Warehouse mounts every registered library when serving `--catalog-config`; +- `--catalog-config` authenticates each request and mounts only granted datasets in isolated workers; - Snapshot refresh does not mutate an in-flight Snapshot; - RustFS Warehouse behavior matches local Datasets for the covered paths. diff --git a/docs/src/en/pchronicle/reference/cli.md b/docs/src/en/pchronicle/reference/cli.md index cfaef8ec2..693249542 100644 --- a/docs/src/en/pchronicle/reference/cli.md +++ b/docs/src/en/pchronicle/reference/cli.md @@ -334,13 +334,27 @@ Every listener must use a loopback address. A bare single Dataset is mounted as needed. Control requires a mount named `default`. Repeatable `--home-link TEXT=PATH` adds homepage nav capsules beside Warehouse. `PATH` must be a same-origin relative path such as `/plugins`. -`--catalog-config FILE` mounts every `[datasets.*]` library in the Directory -file into Warehouse and enables `catalog://` locators. It conflicts with -positional Dataset mounts. Pair Directory clients with +`--catalog-config FILE` authenticates every data API request and dispatches it to +an isolated exec worker containing only authorized mounts. It also enables +`catalog://` locators, and conflicts with positional mounts, Gateway, and Control. Pair Directory clients with `dataset pin NAME catalog://127.0.0.1:PORT --ak --sk`. `pchronicle serve catalog dataset add|remove|list` and `issue|grant|revoke` rewrite that file and do not start HTTP; `issue` prints the user secret once. Restart serve after changing libraries, users, or grants. +The pool allows at most 8 workers and 32 admitted requests, with serial execution +per worker, a 60-second execution/queue timeout and a 10-second body-read timeout. +Overload returns 503; timeout or IPC failure discards the worker. Worker and disk +cache identity includes the user, grants and backend credential version. Caches +live under `PCHRONICLE_CACHE_DIR/workers/` or the system pchronicle cache directory. +Workers receive backend keys over private IPC before starting runtime threads; +they do not inherit AWS environment, profiles or the server's home configuration. +Login AK/SK authenticate the user; dataset AK/SK authenticate storage access. +Datasets with different S3 endpoints or credentials require explicit `dataset` +selection; cross-credential queries are currently rejected. Health, UI config +and static pages remain public; data APIs require authentication even for a +catalog with no users. Processes still use the server's OS identity: this is not +an OS privilege drop or a filesystem sandbox. + `catalog` is a reserved `serve` subcommand; mount a path of that name as `./catalog`. See [RFC-0013](../../rfcs/0013-pchronicle-warehouse-catalog.md) and diff --git a/docs/src/en/rfcs/0016-pchronicle-catalog-resolution-and-cache.md b/docs/src/en/rfcs/0016-pchronicle-catalog-resolution-and-cache.md new file mode 100644 index 000000000..43f6797ea --- /dev/null +++ b/docs/src/en/rfcs/0016-pchronicle-catalog-resolution-and-cache.md @@ -0,0 +1,135 @@ +# RFC-0016: pChronicle Catalog Resolution and Manifest Cache + +| Field | Value | +|---|---| +| **Status** | Accepted | +| **Date** | 2026-09-13 | +| **Component** | `persisting-pchronicle`, `pchronicle` CLI, Warehouse Datasets UI | +| **Related** | [RFC-0013 path Directory](0013-pchronicle-warehouse-catalog.md) · [RFC-0015 `chronicle.manifest`](0015-chronicle-manifest.md) | + +## Summary + +pChronicle has one catalog boundary for locating and describing data. A +`DatasetLocation` identifies a physical location, a `Dataset` identifies one +atomic dataset, and a `DatasetMount` identifies a path whose descendants form +an aggregate query space. Resolution and manifest caching live in the core +`persisting-pchronicle` crate and are shared by the CLI, Web UI, and query +engine. + +The cache is an acceleration layer for discovery and navigation. It MUST NOT +change the correctness contract of queries or CLI operations. + +## Terminology and model + +- **Location**: a parsed local or object-store URI. It owns one-level listing + and manifest reads; it does not decide how a caller uses the result. +- **Dataset**: one atomic, queryable leaf. Its membership and revision are + determined by the authoritative dataset layout and its `chronicle.manifest` + when present. +- **Dataset mount**: a named path that may contain many datasets. It is a + navigational and query boundary, not an additional physical dataset. +- **Manifest listing**: one cached `list(prefix)` observation, including child + kinds and optional record statistics. + +The model is deliberately separate from DataFusion. DataFusion receives a +resolved dataset or mount and applies mount-specific pruning; it does not own +filesystem traversal or the UI cache. + +## Resolver contract + +All callers use the shared resolver and select an explicit read mode: + +| Mode | Purpose | Authoritative I/O | +|---|---|---| +| `Cached` | Render an existing UI view immediately | Never | +| `RefreshIfMissing` | Fill a missing cache entry | Only on a miss | +| `Fresh` | CLI, mutations, and correctness-sensitive operations | Always | + +The resolver MUST normalize prefixes, reject traversal components, preserve +the distinction between an atomic Dataset and a Dataset mount, and keep URI +identity in cache keys. A cached result MAY be stale or incomplete; a fresh +result MUST reflect the authoritative location at the time of resolution. + +CLI catalog calculations and query planning use `Fresh`. The Web UI uses +`Cached` first and refreshes in the background. UI cache data MUST NOT be used +to answer a user query, establish source membership, or select a revision. + +## Manifest cache + +`ManifestCache` is the single owner of manifest observations. Callers do not +open `chronicle.manifest` or maintain a second catalog cache themselves. + +Each entry is keyed by a stable mount identity and normalized prefix: + +```text +dataset-name NUL location-fingerprint [NUL prefix] +``` + +The cache has an in-memory read path backed by a persistent Lance table under +the pChronicle cache directory. Opening a corrupt, incompatible, or +unreadable cache MUST rebuild it and continue with an empty cache. Cache +errors MUST NOT make `serve` fail when the authoritative location is still +available. + +The cache exposes three refresh paths: + +1. A foreground tree request refreshes only the requested level when its + cached view is missing or explicitly stale. +2. A periodic catalog refresh runs every 30 seconds and performs a breadth + first walk of each mount. This publishes shallow levels before deeper + levels and bounds the frontier. +3. A successful refresh updates memory and persistent storage atomically from + the caller's point of view. A failed refresh leaves the previous view + available and records the error for the UI. + +At most one manifest refresh is active per cache instance. This prevents +duplicate timers and concurrent scans from multiplying local or object-store +I/O. The UI browse worker also has a bounded process-level I/O gate. + +## Aggregation semantics + +For a mount or prefix, cached summaries include Dataset leaves observed at that +prefix and all cached descendant prefixes. They report dataset count and +trajectory count when those values are present in the manifest. They are +best-effort UI metadata and can temporarily be zero while a cold cache is +warming. + +Directory cards may display the summary for their child prefix. A directory +without a cached descendant observation remains navigable and displays a +neutral directory label until background refresh discovers its contents. + +Refreshes also remove descendants that disappear from a successfully listed +parent. Failed or unavailable parent reads MUST NOT erase an existing cached +view. + +## Query behavior for mounts + +The query engine MAY accept either an atomic Dataset or a Dataset mount. For a +mount it MUST push path and dataset predicates down before constructing +cross-dataset joins, prune datasets that cannot satisfy the predicate, and +avoid materializing unrelated descendants. These optimizations are part of +query execution, not manifest-cache correctness, and must remain valid when +the cache is stale. + +## Consistency and recovery + +The cache is rebuildable. It is safe to delete the local cache and restart; +the next foreground request or periodic walk repopulates it. Persistent cache +writes are best effort for UI use, while fresh CLI and query paths surface +authoritative read errors. + +This design intentionally does not promise a single global snapshot across +multiple object-store locations. Each manifest listing records its own +observation time; a UI response identifies itself as best effort and exposes +refreshing, stale, and last-error state. + +## Rejected alternatives + +- Re-scanning every manifest for every UI request: predictable worst-case + latency and excessive object-store I/O. +- A CLI-only cache: duplicates resolution rules and allows UI, CLI, and query + semantics to drift. +- Treating a mount as one physical Dataset: hides membership boundaries and + encourages unbounded joins. +- Using cached membership for accurate queries: stale caches can omit newly + created datasets or retain removed ones. diff --git a/docs/src/en/rfcs/index.md b/docs/src/en/rfcs/index.md index 224049dbe..58b2d422e 100644 --- a/docs/src/en/rfcs/index.md +++ b/docs/src/en/rfcs/index.md @@ -26,3 +26,4 @@ each product's Reference and Guides. | [0013](0013-pchronicle-warehouse-catalog.md) | pChronicle path Directory | Proposed | | [0014](0014-compact-jsonl.md) | Compact JSONL Lance storage format (`compact-jsonl/v1`) | Accepted | | [0015](0015-chronicle-manifest.md) | `chronicle.manifest` Dataset sidecar (TOML) | Proposed | +| [0016](0016-pchronicle-catalog-resolution-and-cache.md) | pChronicle catalog resolution and manifest cache | Accepted | diff --git a/docs/src/zh/pchronicle/design/architecture.md b/docs/src/zh/pchronicle/design/architecture.md index 85a68ce0a..089850de2 100644 --- a/docs/src/zh/pchronicle/design/architecture.md +++ b/docs/src/zh/pchronicle/design/architecture.md @@ -130,10 +130,10 @@ Server 静态挂载命名 path。Refresh 先完整构造新 Snapshot,再切换 Dataset table 先按 Source 裁剪,再打开命中的固定 version;cache 和 routing index 与 Snapshot generation 绑定。 -使用 `--catalog-config` 时,Warehouse 会把 Directory ACL 文件中的全部 -`[datasets.*]` library 挂进数据面(与位置参数挂载等价),并同时提供 -`catalog://` 列表/换票路由。文件中的 S3 endpoint、region 与后端密钥在打开存储前 -写入进程环境。CLI 换票后打开票里的 `uri`(一条 path)并注入存储钥。这是 path 上的 +使用 `--catalog-config` 时,父进程逐请求认证,并调度到有上限的 exec worker 池。 +每个 worker 的用户、授权及后端凭证范围固定,缓存独立;父进程不挂载数据集。 +后端凭证经私有 IPC 在 runtime 启动前传入。worker 仍使用服务端 OS 身份,进程隔离 +不等于文件系统沙箱。同时保留 `catalog://` 列表/换票路由。CLI 换票后打开票里的 `uri`(一条 path)并注入存储钥。这是 path 上的 平台寻址,不是新的 Dataset 种类。 Web 与 API 是同一读取模型的 consumer,不形成新事实源。未知 API route 保持 error,不进入 diff --git a/docs/src/zh/pchronicle/guides/serve.md b/docs/src/zh/pchronicle/guides/serve.md index a04de435a..ead680d54 100644 --- a/docs/src/zh/pchronicle/guides/serve.md +++ b/docs/src/zh/pchronicle/guides/serve.md @@ -65,6 +65,7 @@ pchronicle serve --catalog-config catalog.toml --listen 127.0.0.1:8081 `serve catalog dataset add|remove|list` 改写 libraries,不启动 HTTP。 `serve catalog issue` 写入一个无授权用户,并把 sk 只打印到这次 stdout; `grant` / `revoke` 改该用户可打开的 library 名称。改文件后必须重启 serve。 +NAME 使用 `*` 可以为当前所有用户授予或撤销数据集。新建用户不会自动继承过去的通配授权,创建后请重新执行命令。 `pchronicle serve --catalog-config` 会把文件中的 **全部** library 挂进 Warehouse (与位置参数挂载等价),并启用 `catalog://` 换票路由。不要与位置参数 Dataset diff --git a/docs/src/zh/pchronicle/guides/ui.md b/docs/src/zh/pchronicle/guides/ui.md index 683b00891..35b5518ba 100644 --- a/docs/src/zh/pchronicle/guides/ui.md +++ b/docs/src/zh/pchronicle/guides/ui.md @@ -129,8 +129,8 @@ Storage 是高级诊断页,不是日常浏览 Run 的必经步骤。左侧按 清除该站点的浏览器数据也会清除这份设置。Assistant 标记为 **Read-only · selected run data**, 用于解释当前上下文,不会改写 Dataset。 -使用 `pchronicle serve --catalog-config` 时,ACL 文件中的全部 library 已挂载供本机浏览。 -若需要 Directory 用户鉴权流程,从左侧 **Keys** 填写 access key 和 secret key。它们保存在 +使用 `pchronicle serve --catalog-config` 时,数据 API 需要认证,仅展示用户获准的数据集。 +从左侧 **Keys** 填写 access key 和 secret key 后访问数据。它们保存在 `localStorage`,并作为 `x-pchronicle-access-key` / `x-pchronicle-secret-key` 发给当前 pChronicle 服务端。它们决定浏览器可打开哪些 Directory path,不是对象存储后端密钥。 diff --git a/docs/src/zh/pchronicle/reference/cases-platform.md b/docs/src/zh/pchronicle/reference/cases-platform.md index 26408762b..e029bb8f8 100644 --- a/docs/src/zh/pchronicle/reference/cases-platform.md +++ b/docs/src/zh/pchronicle/reference/cases-platform.md @@ -92,6 +92,6 @@ export PCHRONICLE_RUSTFS_BUCKET=pchronicle-cases - ACL 可从空文件开始构建; - 用户、dataset 与 grants 修改是确定性的; - 后端对象存储密钥留在 catalog 文件 / ticket 路径,不出现在 `dataset list`; -- `--catalog-config` serve 会挂载全部已登记 library; +- `--catalog-config` serve 逐请求认证,独立 worker 仅挂载该用户获准的数据集; - Snapshot refresh 不改动进行中查询的 Snapshot; - 覆盖路径上 RustFS Warehouse 行为与本地 Dataset 一致。 diff --git a/docs/src/zh/pchronicle/reference/cli.md b/docs/src/zh/pchronicle/reference/cli.md index 43b4b58cf..f9cc6e8f0 100644 --- a/docs/src/zh/pchronicle/reference/cli.md +++ b/docs/src/zh/pchronicle/reference/cli.md @@ -439,10 +439,20 @@ pchronicle serve \ `NAME=DATASET` mount;Control 模式要求名为 `default` 的 mount。可重复的 `--home-link TEXT=PATH` 会在首页 Warehouse 旁增加胶囊;`PATH` 必须是同源相对路径。 `--catalog-config FILE` -会把文件中全部 `[datasets.*]` 挂进 Warehouse,并启用 `catalog://` locator;不能与位置参数 -Dataset 同时使用。配合 `dataset pin NAME catalog://127.0.0.1:PORT --ak --sk`。 +启用逐请求 AK/SK 认证及 `catalog://` locator。父进程只监听、认证和调度,已授权数据集由 +独立 exec worker 读取;不能与位置参数 Dataset、Gateway 或 Control 同时使用。配合 `dataset pin NAME catalog://127.0.0.1:PORT --ak --sk`。 `pchronicle serve catalog dataset add|remove|list` 与 `issue|grant|revoke` 只改该文件、 不启动 HTTP;`issue` 把用户 sk 只打印一次。改 library、用户或授权后必须重启 serve。 +worker 池最多 8 个进程、32 个正在处理或排队的请求;每个 worker 串行处理请求, +计算等待上限为 60 秒,请求体读取上限为 10 秒。过载返回 503,超时或 IPC 失败会淘汰进程。 +同一用户、授权范围和后端凭证版本复用 worker 及独立磁盘缓存;修改授权后重启生效。 +缓存位于 `PCHRONICLE_CACHE_DIR/workers/`,未设置时使用系统 pchronicle 缓存目录。 +子进程不继承父进程的 AWS 环境、profile 或用户主目录配置;登录 AK/SK 用于认证, +数据集的 AK/SK 经私有 IPC 传入并用于对象存储访问。不同 S3 endpoint/凭证的数据集必须 +分别通过 `dataset` 参数选择,暂不支持这类跨凭证查询。健康检查、UI 配置和静态页面公开, +数据 API 需要认证。没有用户的 catalog 仍可管理,但数据 API 不接受匿名访问。 +这是进程与凭证上下文隔离,不是 OS 用户降权或文件系统沙箱;本地路径仍使用服务进程的 OS 身份。 + `catalog` 是 `serve` 的保留子命令,挂载同名路径请用 `./catalog`。见 [RFC-0013](../../rfcs/0013-pchronicle-warehouse-catalog.md) 与 [RFC-0015](../../rfcs/0015-chronicle-manifest.md)。无需配置的 `--gateway` diff --git a/docs/src/zh/rfcs/0016-pchronicle-catalog-resolution-and-cache.md b/docs/src/zh/rfcs/0016-pchronicle-catalog-resolution-and-cache.md new file mode 100644 index 000000000..43f6797ea --- /dev/null +++ b/docs/src/zh/rfcs/0016-pchronicle-catalog-resolution-and-cache.md @@ -0,0 +1,135 @@ +# RFC-0016: pChronicle Catalog Resolution and Manifest Cache + +| Field | Value | +|---|---| +| **Status** | Accepted | +| **Date** | 2026-09-13 | +| **Component** | `persisting-pchronicle`, `pchronicle` CLI, Warehouse Datasets UI | +| **Related** | [RFC-0013 path Directory](0013-pchronicle-warehouse-catalog.md) · [RFC-0015 `chronicle.manifest`](0015-chronicle-manifest.md) | + +## Summary + +pChronicle has one catalog boundary for locating and describing data. A +`DatasetLocation` identifies a physical location, a `Dataset` identifies one +atomic dataset, and a `DatasetMount` identifies a path whose descendants form +an aggregate query space. Resolution and manifest caching live in the core +`persisting-pchronicle` crate and are shared by the CLI, Web UI, and query +engine. + +The cache is an acceleration layer for discovery and navigation. It MUST NOT +change the correctness contract of queries or CLI operations. + +## Terminology and model + +- **Location**: a parsed local or object-store URI. It owns one-level listing + and manifest reads; it does not decide how a caller uses the result. +- **Dataset**: one atomic, queryable leaf. Its membership and revision are + determined by the authoritative dataset layout and its `chronicle.manifest` + when present. +- **Dataset mount**: a named path that may contain many datasets. It is a + navigational and query boundary, not an additional physical dataset. +- **Manifest listing**: one cached `list(prefix)` observation, including child + kinds and optional record statistics. + +The model is deliberately separate from DataFusion. DataFusion receives a +resolved dataset or mount and applies mount-specific pruning; it does not own +filesystem traversal or the UI cache. + +## Resolver contract + +All callers use the shared resolver and select an explicit read mode: + +| Mode | Purpose | Authoritative I/O | +|---|---|---| +| `Cached` | Render an existing UI view immediately | Never | +| `RefreshIfMissing` | Fill a missing cache entry | Only on a miss | +| `Fresh` | CLI, mutations, and correctness-sensitive operations | Always | + +The resolver MUST normalize prefixes, reject traversal components, preserve +the distinction between an atomic Dataset and a Dataset mount, and keep URI +identity in cache keys. A cached result MAY be stale or incomplete; a fresh +result MUST reflect the authoritative location at the time of resolution. + +CLI catalog calculations and query planning use `Fresh`. The Web UI uses +`Cached` first and refreshes in the background. UI cache data MUST NOT be used +to answer a user query, establish source membership, or select a revision. + +## Manifest cache + +`ManifestCache` is the single owner of manifest observations. Callers do not +open `chronicle.manifest` or maintain a second catalog cache themselves. + +Each entry is keyed by a stable mount identity and normalized prefix: + +```text +dataset-name NUL location-fingerprint [NUL prefix] +``` + +The cache has an in-memory read path backed by a persistent Lance table under +the pChronicle cache directory. Opening a corrupt, incompatible, or +unreadable cache MUST rebuild it and continue with an empty cache. Cache +errors MUST NOT make `serve` fail when the authoritative location is still +available. + +The cache exposes three refresh paths: + +1. A foreground tree request refreshes only the requested level when its + cached view is missing or explicitly stale. +2. A periodic catalog refresh runs every 30 seconds and performs a breadth + first walk of each mount. This publishes shallow levels before deeper + levels and bounds the frontier. +3. A successful refresh updates memory and persistent storage atomically from + the caller's point of view. A failed refresh leaves the previous view + available and records the error for the UI. + +At most one manifest refresh is active per cache instance. This prevents +duplicate timers and concurrent scans from multiplying local or object-store +I/O. The UI browse worker also has a bounded process-level I/O gate. + +## Aggregation semantics + +For a mount or prefix, cached summaries include Dataset leaves observed at that +prefix and all cached descendant prefixes. They report dataset count and +trajectory count when those values are present in the manifest. They are +best-effort UI metadata and can temporarily be zero while a cold cache is +warming. + +Directory cards may display the summary for their child prefix. A directory +without a cached descendant observation remains navigable and displays a +neutral directory label until background refresh discovers its contents. + +Refreshes also remove descendants that disappear from a successfully listed +parent. Failed or unavailable parent reads MUST NOT erase an existing cached +view. + +## Query behavior for mounts + +The query engine MAY accept either an atomic Dataset or a Dataset mount. For a +mount it MUST push path and dataset predicates down before constructing +cross-dataset joins, prune datasets that cannot satisfy the predicate, and +avoid materializing unrelated descendants. These optimizations are part of +query execution, not manifest-cache correctness, and must remain valid when +the cache is stale. + +## Consistency and recovery + +The cache is rebuildable. It is safe to delete the local cache and restart; +the next foreground request or periodic walk repopulates it. Persistent cache +writes are best effort for UI use, while fresh CLI and query paths surface +authoritative read errors. + +This design intentionally does not promise a single global snapshot across +multiple object-store locations. Each manifest listing records its own +observation time; a UI response identifies itself as best effort and exposes +refreshing, stale, and last-error state. + +## Rejected alternatives + +- Re-scanning every manifest for every UI request: predictable worst-case + latency and excessive object-store I/O. +- A CLI-only cache: duplicates resolution rules and allows UI, CLI, and query + semantics to drift. +- Treating a mount as one physical Dataset: hides membership boundaries and + encourages unbounded joins. +- Using cached membership for accurate queries: stale caches can omit newly + created datasets or retain removed ones. diff --git a/docs/src/zh/rfcs/index.md b/docs/src/zh/rfcs/index.md index b35b8b968..73de532fd 100644 --- a/docs/src/zh/rfcs/index.md +++ b/docs/src/zh/rfcs/index.md @@ -22,3 +22,4 @@ RFC 正文保持英文,因为它们是历史决策快照。翻译会产生两 | [0013](0013-pchronicle-warehouse-catalog.md) | pChronicle path Directory | Proposed | | [0014](0014-compact-jsonl.md) | Compact JSONL Lance 存储格式(`compact-jsonl/v1`) | Accepted | | [0015](0015-chronicle-manifest.md) | `chronicle.manifest` Dataset sidecar(TOML) | Proposed | +| [0016](0016-pchronicle-catalog-resolution-and-cache.md) | pChronicle catalog resolve 与 manifest cache | Accepted | diff --git a/pchronicle-web/assets/app.css b/pchronicle-web/assets/app.css index a5b6a7811..518d221ab 100644 --- a/pchronicle-web/assets/app.css +++ b/pchronicle-web/assets/app.css @@ -1 +1,2 @@ :root{font-family:Inter,ui-sans-serif,-apple-system,BlinkMacSystemFont,"Segoe UI",sans-serif;color:#111827;background:#f6f7f9;font-synthesis:none;--blue:#2563eb;--slate:#0b1220;--border:#e2e5ea;--muted:#667085}*{box-sizing:border-box}html,body,#main{height:100%;margin:0}body{overflow:hidden}button,input,select,textarea{font:inherit}button,a{outline:none}.app-shell{height:100vh;display:grid;grid-template-columns:56px 232px minmax(0,1fr);background:#f6f7f9}.skip-link{position:fixed;left:12px;top:-60px;z-index:100;background:#fff;color:#1d4ed8;padding:9px 12px;border-radius:7px;box-shadow:0 8px 20px #0003}.skip-link:focus{top:12px}.rail{display:flex;flex-direction:column;align-items:center;gap:7px;padding:12px 7px;background:linear-gradient(180deg,#101b2f,#07101f);border-right:1px solid #263247;color:#cbd5e1}.brand-mark{width:38px;height:38px;display:grid;place-items:center;margin-bottom:10px;border:1px solid #3b82f680;border-radius:11px;background:linear-gradient(145deg,#2563eb,#1d4ed8);font-size:13px;font-weight:800;letter-spacing:-.04em;color:#fff;box-shadow:0 8px 20px #1d4ed84d;padding:0;cursor:pointer}.rail-button{width:42px;min-height:50px;display:flex;flex-direction:column;align-items:center;justify-content:center;gap:3px;border:0;border-radius:9px;background:transparent;color:#94a3b8;font-size:10px;cursor:pointer}.rail-button:hover,.rail-button:focus-visible{background:#ffffff10;color:#e2e8f0}.rail-button.active{background:#2563eb22;color:#93c5fd;box-shadow:inset 0 0 0 1px #3b82f655}.rail-icon{font-size:19px;line-height:1}.rail-spacer{flex:1}.rail-status{display:flex;flex-direction:column;align-items:center;gap:4px;color:#64748b;font-size:9px}.live-dot{display:inline-block;width:7px;height:7px;border-radius:50%;background:#22c55e;box-shadow:0 0 0 3px #22c55e1c}.run-sidebar{min-width:0;display:flex;flex-direction:column;background:#0d1726;color:#e5edf7;border-right:1px solid #1f2d40}.sidebar-heading{height:76px;display:flex;align-items:center;justify-content:space-between;padding:12px 14px;border-bottom:1px solid #1f2d40}.eyebrow{margin:0 0 4px;color:#60a5fa;font-size:10px;font-weight:700;text-transform:uppercase;letter-spacing:.1em}.sidebar-heading h1{margin:0;font-size:15px}.icon-button{display:grid;place-items:center;min-width:30px;height:30px;border:1px solid transparent;border-radius:7px;background:transparent;color:inherit;cursor:pointer}.icon-button:hover,.icon-button:focus-visible{background:#ffffff0d;border-color:#ffffff1a}.search-field{display:flex;align-items:center;gap:7px;margin:12px;padding:8px 9px;border:1px solid #2b3a4f;border-radius:8px;background:#101e30;color:#8291a5}.search-field:focus-within{border-color:#3b82f6;box-shadow:0 0 0 3px #2563eb20}.search-field input{min-width:0;flex:1;border:0;outline:0;background:transparent;color:#e5edf7;font-size:12px}.search-field input::placeholder{color:#68788d}.search-field kbd{padding:1px 5px;border:1px solid #34455c;border-radius:4px;font-size:10px}.run-count{padding:0 14px 7px;color:#718198;font-size:10px;text-transform:uppercase;letter-spacing:.08em}.run-list{min-height:0;flex:1;overflow:auto;padding:0 8px 14px}.run-item{width:100%;display:block;margin-bottom:4px;padding:10px;border:1px solid transparent;border-radius:8px;background:transparent;color:#cbd5e1;text-align:left;cursor:pointer}.run-item:hover{background:#142236}.run-item.selected{border-color:#2f67b6;background:#162a46;box-shadow:inset 3px 0 #3b82f6}.run-item-top,.run-meta{display:flex;align-items:center;justify-content:space-between;gap:8px}.run-item-top strong{overflow:hidden;text-overflow:ellipsis;font-size:12px;white-space:nowrap}.run-session{margin:5px 0;overflow:hidden;color:#9fb0c4;font:11px ui-monospace,SFMono-Regular,Menlo,monospace;text-overflow:ellipsis;white-space:nowrap}.run-meta{color:#718198;font-size:10px}.warning-text{color:#fbbf24}.run-skeleton{height:64px;margin:4px 0;border-radius:8px;background:linear-gradient(90deg,#132033,#1d2c41,#132033);background-size:200% 100%;animation:shimmer 1.4s infinite}@keyframes shimmer{to{background-position:-200% 0}}.workspace{min-width:0;min-height:0;display:flex;flex-direction:column;overflow:hidden}.workspace-header{min-height:76px;display:flex;align-items:center;justify-content:space-between;gap:20px;padding:12px 20px;border-bottom:1px solid var(--border);background:#fff}.title-block{min-width:0}.breadcrumb{color:#667085;font-size:11px}.title-block h2{margin:3px 0 4px;overflow:hidden;font-size:18px;line-height:1.2;text-overflow:ellipsis;white-space:nowrap}.header-meta,.header-actions{display:flex;align-items:center;gap:8px;color:#667085;font-size:11px}.header-meta code{max-width:240px;overflow:hidden;color:#475467;text-overflow:ellipsis}.header-actions{flex-shrink:0}.button{display:inline-flex;align-items:center;gap:7px;padding:7px 10px;border:1px solid #d0d5dd;border-radius:7px;background:#fff;color:#344054;text-decoration:none;font-size:11px;font-weight:600;cursor:pointer}.button:hover,.button:focus-visible{border-color:#98a2b3;background:#f9fafb}.button.primary{border-color:#2563eb;background:#2563eb;color:#fff}.button.danger{border-color:#fecaca;background:#fff5f5;color:#b42318}.button.active-follow{border-color:#bbf7d0;background:#f0fdf4;color:#166534}.status-pill{display:inline-flex;align-items:center;gap:5px;padding:2px 6px;border:1px solid #d0d5dd;border-radius:999px;background:#fff;color:#475467;font-size:9px;font-weight:700;text-transform:uppercase}.status-pill.good{border-color:#bbf7d0;background:#f0fdf4;color:#15803d}.status-dot,.kind-dot{width:5px;height:5px;border-radius:50%;background:currentColor}.notice{display:flex;align-items:center;gap:10px;margin:12px 20px 0;padding:9px 12px;border:1px solid;border-radius:8px;font-size:11px}.error-notice{border-color:#fecaca;background:#fff5f5;color:#991b1b}.new-events{position:absolute;z-index:20;left:50%;top:86px;transform:translateX(-50%);padding:7px 12px;border:1px solid #93c5fd;border-radius:999px;background:#eff6ff;color:#1d4ed8;font-size:11px;font-weight:700;box-shadow:0 5px 15px #1d4ed822;cursor:pointer}.metric-strip{display:grid;grid-template-columns:repeat(4,minmax(0,1fr));margin:14px 20px 12px;border:1px solid var(--border);border-radius:9px;background:#fff}.metric{min-width:0;display:grid;grid-template-columns:1fr auto;gap:2px 10px;padding:10px 14px;border-right:1px solid #eceef1}.metric:last-child{border-right:0}.metric>span{color:#667085;font-size:10px;font-weight:600;text-transform:uppercase;letter-spacing:.05em}.metric strong{grid-row:1/3;grid-column:2;font:600 20px ui-monospace,SFMono-Regular,Menlo,monospace;color:#101828}.metric small{overflow:hidden;color:#98a2b3;font-size:10px;text-overflow:ellipsis;white-space:nowrap}.evidence-layout{min-height:0;flex:1;display:grid;grid-template-columns:minmax(0,1fr) 360px;gap:12px;padding:0 20px 18px}.evidence-layout.inspector-hidden{grid-template-columns:minmax(0,1fr)}.evidence-surface,.inspector{min-width:0;min-height:0;display:flex;flex-direction:column;border:1px solid var(--border);border-radius:10px;background:#fff;overflow:hidden}.surface-toolbar{min-height:54px;display:flex;align-items:center;justify-content:space-between;gap:12px;padding:8px 12px;border-bottom:1px solid #eceef1}.surface-title{font-size:13px;font-weight:700}.surface-subtitle{margin-top:2px;color:#667085;font-size:10px}.filters{display:flex;gap:6px}.filters input,.filters select{height:30px;border:1px solid #d0d5dd;border-radius:6px;background:#fff;color:#344054;font-size:11px}.filters input{width:190px;padding:0 9px}.filters select{padding:0 26px 0 8px}.trajectory-scroll{min-height:0;flex:1;overflow:auto;padding:10px 14px 26px;scrollbar-gutter:stable}.time-ruler{display:flex;align-items:center;gap:9px;margin:0 0 8px 34px;color:#98a2b3;font:9px ui-monospace,SFMono-Regular,Menlo,monospace;text-transform:uppercase}.ruler-line{height:1px;flex:1;background:linear-gradient(90deg,#d0d5dd,#e5e7eb)}.turn-row{display:grid;grid-template-columns:24px minmax(0,1fr);gap:10px;cursor:pointer}.turn-row:focus-visible .turn-card,.turn-row.selected .turn-card{border-color:#93c5fd;box-shadow:0 0 0 3px #2563eb14}.turn-axis{display:flex;flex-direction:column;align-items:center}.turn-dot{z-index:1;width:10px;height:10px;margin-top:14px;border:2px solid #fff;border-radius:50%;background:#94a3b8;box-shadow:0 0 0 1px #cbd5e1}.turn-dot.user{background:#2563eb}.turn-dot.agent{background:#10b981}.turn-dot.system{background:#f59e0b}.turn-line{width:1px;min-height:30px;flex:1;background:#d7dce2}.turn-card{margin-bottom:9px;border:1px solid #e4e7ec;border-radius:8px;background:#fff;overflow:hidden;transition:border-color .12s,box-shadow .12s}.turn-card:hover{border-color:#cbd5e1}.turn-header,.turn-footer{display:flex;align-items:center;justify-content:space-between;gap:10px;padding:7px 9px;background:#fafbfc;color:#667085;font-size:9px}.turn-header{border-bottom:1px solid #f0f1f3}.turn-identity,.turn-timing{display:flex;align-items:center;gap:7px}.turn-timing time{max-width:190px;overflow:hidden;text-overflow:ellipsis;white-space:nowrap}.role-badge{padding:2px 6px;border-radius:4px;background:#f2f4f7;color:#475467;font-size:9px;font-weight:800;text-transform:uppercase}.role-badge.user{background:#eff6ff;color:#1d4ed8}.role-badge.agent{background:#ecfdf5;color:#047857}.role-badge.system{background:#fffbeb;color:#b45309}.turn-content{max-height:360px;margin:0;padding:10px 12px;overflow:auto;background:#fff;color:#27364a;font:11px/1.55 ui-monospace,SFMono-Regular,Menlo,monospace;white-space:pre-wrap;word-break:break-word}.turn-footer{justify-content:flex-start;border-top:1px solid #f0f1f3}.reasoning{margin:0 10px 8px;border:1px solid #e4e7ec;border-radius:6px;color:#475467;font-size:10px}.reasoning summary{padding:6px 8px;cursor:pointer}.reasoning pre,.tool-call pre{margin:0;padding:8px;border-top:1px solid #e4e7ec;overflow:auto;white-space:pre-wrap}.tool-call{margin:0 10px 8px;border:1px solid #bfdbfe;border-radius:7px;background:#f8fbff;font-size:10px}.tool-call>div{display:flex;justify-content:space-between;padding:7px 8px;color:#1e40af}.inspector-header{min-height:54px;display:flex;align-items:center;justify-content:space-between;padding:8px 12px;border-bottom:1px solid #eceef1}.inspector-header h3{margin:0;font-size:13px}.inspector-tabs{display:flex;padding:0 10px;border-bottom:1px solid #eceef1}.inspector-tabs button{padding:9px 7px;border:0;border-bottom:2px solid transparent;background:transparent;color:#667085;font-size:10px;cursor:pointer}.inspector-tabs button.active{border-color:#2563eb;color:#1d4ed8;font-weight:700}.inspector-body{min-height:0;flex:1;overflow:auto;padding:10px}.inspector-field{display:grid;grid-template-columns:88px minmax(0,1fr);gap:8px;padding:7px 0;border-bottom:1px solid #f0f1f3;font-size:10px}.inspector-field span,.inspector-code>span{color:#667085}.inspector-field code{overflow:hidden;color:#344054;text-overflow:ellipsis;white-space:nowrap}.inspector-code{margin-top:12px;font-size:10px}.inspector-code pre,.raw-event pre{margin:5px 0 0;padding:9px;border:1px solid #e4e7ec;border-radius:6px;background:#f8fafc;color:#344054;font:9px/1.5 ui-monospace,SFMono-Regular,Menlo,monospace;white-space:pre-wrap;word-break:break-word}.raw-event{margin-bottom:8px;border:1px solid #e4e7ec;border-radius:7px}.raw-event summary{display:flex;align-items:center;gap:7px;padding:8px;color:#475467;font-size:9px;cursor:pointer}.raw-event summary span:last-child{margin-left:auto;color:#98a2b3}.raw-event pre{margin:0;border:0;border-top:1px solid #e4e7ec;border-radius:0}.inspector-empty{padding:20px;color:#667085;font-size:11px;line-height:1.6}.loading-panel,.empty-state{display:flex;min-height:170px;flex-direction:column;align-items:center;justify-content:center;color:#667085;text-align:center}.empty-state strong{margin-top:7px;color:#344054;font-size:12px}.empty-state p{max-width:260px;margin:5px 0;font-size:10px;line-height:1.5}.empty-icon{font-size:22px;color:#98a2b3}.spinner{width:16px;height:16px;margin-bottom:8px;border:2px solid #dbeafe;border-top-color:#2563eb;border-radius:50%;animation:spin .8s linear infinite}@keyframes spin{to{transform:rotate(360deg)}}.welcome-state{max-width:580px;margin:auto;padding:50px;text-align:center}.welcome-orbit{position:relative;width:88px;height:88px;display:grid;place-items:center;margin:0 auto 20px;border:1px solid #bfdbfe;border-radius:50%;background:#eff6ff;color:#1d4ed8;font-weight:800;box-shadow:0 0 0 16px #eff6ff80}.orbit-dot{position:absolute;top:7px;right:12px;width:8px;height:8px;border-radius:50%;background:#10b981;box-shadow:0 0 0 4px #d1fae5}.welcome-state h2{margin:5px 0 8px;font-size:24px}.welcome-state>p:not(.eyebrow){margin:0;color:#667085;font-size:13px;line-height:1.65}.welcome-keys{display:flex;justify-content:center;gap:20px;margin-top:24px;color:#667085;font-size:10px}.welcome-keys kbd{margin-right:5px;padding:3px 6px;border:1px solid #d0d5dd;border-radius:5px;background:#fff;color:#344054}.tools-workspace{min-height:0;display:flex;flex:1;flex-direction:column}.tools-grid{min-height:0;display:grid;grid-template-columns:220px minmax(0,1fr);gap:14px;flex:1;padding:18px 20px}.tools-nav,.tool-surface{border:1px solid var(--border);border-radius:10px;background:#fff}.tools-nav{padding:7px}.tools-nav button{width:100%;display:flex;flex-direction:column;gap:3px;padding:10px;border:0;border-radius:7px;background:transparent;color:#344054;text-align:left;cursor:pointer}.tools-nav button:hover{background:#f8fafc}.tools-nav button.active{background:#eff6ff;color:#1d4ed8}.tools-nav strong{font-size:11px}.tools-nav span{color:#98a2b3;font-size:9px}.tool-surface{min-width:0;min-height:0;padding:16px;overflow:auto}.tool-heading h3{margin:0;font-size:14px}.tool-heading p{margin:4px 0 14px;color:#667085;font-size:10px}.danger-heading{padding:10px;border:1px solid #fecaca;border-radius:7px;background:#fff8f8}.sql-editor{width:100%;min-height:120px;margin-bottom:9px;padding:11px;border:1px solid #d0d5dd;border-radius:7px;outline:0;background:#101828;color:#d1e9ff;font:11px/1.55 ui-monospace,SFMono-Regular,Menlo,monospace;resize:vertical}.sql-editor:focus{border-color:#3b82f6;box-shadow:0 0 0 3px #2563eb1a}.tool-output{margin-top:16px;border:1px solid #e4e7ec;border-radius:8px;overflow:hidden}.tool-output>div{display:flex;align-items:center;justify-content:space-between;padding:6px 10px;border-bottom:1px solid #e4e7ec;background:#f8fafc;color:#667085;font-size:9px;text-transform:uppercase;letter-spacing:.08em}.tool-output pre{min-height:180px;margin:0;padding:12px;overflow:auto;background:#fff;color:#344054;font:10px/1.55 ui-monospace,SFMono-Regular,Menlo,monospace;white-space:pre-wrap}@media(max-width:1050px){.evidence-layout{grid-template-columns:minmax(520px,1fr) 320px}.metric-strip{grid-template-columns:repeat(2,1fr)}.metric:nth-child(2){border-right:0}.metric:nth-child(-n+2){border-bottom:1px solid #eceef1}.header-actions a{display:none}}@media(prefers-reduced-motion:reduce){*{scroll-behavior:auto!important;animation-duration:.01ms!important;animation-iteration-count:1!important}} +.live-dot.muted{background:#94a3b8;box-shadow:none} diff --git a/pchronicle-web/assets/workbench.css b/pchronicle-web/assets/workbench.css index cd039cb0f..45f22ce1e 100644 --- a/pchronicle-web/assets/workbench.css +++ b/pchronicle-web/assets/workbench.css @@ -149,3 +149,4 @@ .pc2-shell .rail-secondary{padding-top:12px;border-top:1px solid #ffffff14} .pc2-shell .rail-secondary .rail-button{color:#b8cbe4} @media(max-width:850px){.pc2-shell .rail-secondary{padding-top:0;border-top:0}} +.pc2-form select{width:100%;height:42px;padding:0 12px;border:1px solid #cbd5e1;border-radius:9px;background:#fff;color:#344054;font:inherit}.pc2-identity-list{display:flex;flex-direction:column;gap:8px}.pc2-identity-list .button{width:100%;text-align:left}.pc2-settings .button.danger{border-color:#fecaca;color:#b42318;background:#fff}.pc2-settings .button.danger:hover{background:#fff1f0} diff --git a/pchronicle-web/src/api.rs b/pchronicle-web/src/api.rs index 521209b4f..dc50a679f 100644 --- a/pchronicle-web/src/api.rs +++ b/pchronicle-web/src/api.rs @@ -176,6 +176,18 @@ pub async fn explorer_tree(dataset: &str, prefix: &str) -> Result Result { + let url = format!( + "/api/explorer/tree?dataset={}&prefix={}", + urlencoding::encode(dataset), + urlencoding::encode(prefix), + ); + json_checked(Request::get(&url).send().await).await +} + pub async fn run_analysis(run: &RunSummary) -> Result { json_checked( with_catalog_headers(Request::get(&format!("/api/explorer/run?{}", run.query()))) diff --git a/pchronicle-web/src/catalog.rs b/pchronicle-web/src/catalog.rs index bae0bc215..dacbb1d4e 100644 --- a/pchronicle-web/src/catalog.rs +++ b/pchronicle-web/src/catalog.rs @@ -9,6 +9,8 @@ pub fn CatalogExplorer( loading: bool, on_open: EventHandler<(String, String)>, on_runs: EventHandler<(String, String)>, + #[props(default = false)] auth_required: bool, + on_settings: EventHandler, ) -> Element { let dataset = tree .as_ref() @@ -55,6 +57,12 @@ pub fn CatalogExplorer( div { class: "pc-catalog-mosaic", if loading && tree.is_none() { div { class: "pc-catalog-empty", span { class: "spinner" } "Loading datasets…" } + } else if auth_required { + div { class: "pc-catalog-empty", + strong { "Catalog identity required" } + span { "Add an access key and secret key to browse this catalog." } + button { class: "button primary", onclick: on_settings, "Open Keys" } + } } else if tree.as_ref().is_none_or(|tree| tree.children.is_empty() && tree.run_count == 0) { div { class: "pc-catalog-empty", strong { "No datasets" } span { "Add a dataset, then refresh this page." } } } else if tree.as_ref().is_some_and(|tree| tree.children.is_empty()) { @@ -124,8 +132,8 @@ fn CatalogStats(tree: Option) -> Element { let errors = tree.error_sources.unwrap_or(0); rsx! { div { class: "pc-catalog-stats", - div { WorkspaceIcon { name: "folder" } div { span { "Items" } strong { "{tree.children.len()}" } } } - div { WorkspaceIcon { name: "analysis" } div { span { "Trajectories" } strong { "{tree.run_count}" } } } + div { WorkspaceIcon { name: "folder" } div { span { "Datasets" } strong { "{tree.dataset_count.unwrap_or_else(|| tree.children.len())}" } } } + div { WorkspaceIcon { name: "analysis" } div { span { "Trajectories" } strong { "{tree.trajectory_count.unwrap_or(tree.run_count)}" } } } div { WorkspaceIcon { name: "warning" } div { span { "Failed" } strong { "{tree.failed_count}" } } } } if errors > 0 { @@ -168,14 +176,15 @@ fn CatalogFolder( let name = child.name.clone(); let data_type = child.data_type.clone(); let is_dir = kind == "dir"; + let trajectory_count = child.trajectory_count.unwrap_or(child.run_count); let icon = if kind == "file" { "file" } else { "folder" }; rsx! { button { class: "pc-catalog-folder type-{data_type} kind-{kind}", title: if is_dir { format!("{name} · directory") - } else if child.run_count > 0 { - format!("{name} · {} trajectories", child.run_count) + } else if trajectory_count > 0 { + format!("{name} · {trajectory_count} trajectories") } else { name.clone() }, @@ -194,9 +203,15 @@ fn CatalogFolder( span { class: "pc-catalog-folder-type", "{data_type}" } div { class: "pc-catalog-folder-meta", if is_dir { - span { "Directory" } - } else if child.run_count > 0 { - span { "{child.run_count} trajectories" } + if let Some(dataset_count) = child.dataset_count { + span { + "{dataset_count} datasets · {child.trajectory_count.unwrap_or(0)} trajectories" + } + } else { + span { "Directory" } + } + } else if trajectory_count > 0 { + span { "{trajectory_count} trajectories" } } else { span { "Source" } } diff --git a/pchronicle-web/src/catalog_auth.rs b/pchronicle-web/src/catalog_auth.rs index 30ebbb175..2506ca245 100644 --- a/pchronicle-web/src/catalog_auth.rs +++ b/pchronicle-web/src/catalog_auth.rs @@ -1,32 +1,98 @@ use serde::{Deserialize, Serialize}; - const ACCESS_KEY_STORAGE: &str = "pchronicle.catalog.access_key"; const SECRET_KEY_STORAGE: &str = "pchronicle.catalog.secret_key"; +const ACCOUNTS_STORAGE: &str = "pchronicle.catalog.accounts"; +const ACTIVE_STORAGE: &str = "pchronicle.catalog.active"; #[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)] pub struct CatalogAuth { + #[serde(default)] + pub label: String, pub access_key: String, pub secret_key: String, } - impl CatalogAuth { pub fn is_configured(&self) -> bool { !self.access_key.trim().is_empty() && !self.secret_key.trim().is_empty() } } - +pub fn accounts() -> Vec { + let values = read_storage(ACCOUNTS_STORAGE) + .and_then(|v| serde_json::from_str(&v).ok()) + .unwrap_or_else(|| { + let legacy = CatalogAuth { + label: String::new(), + access_key: read_storage(ACCESS_KEY_STORAGE).unwrap_or_default(), + secret_key: read_storage(SECRET_KEY_STORAGE).unwrap_or_default(), + }; + legacy + .is_configured() + .then_some(vec![legacy]) + .unwrap_or_default() + }); + values + .into_iter() + .map(|mut auth: CatalogAuth| { + if auth.label.trim().is_empty() { + auth.label = auth.access_key.clone(); + } + auth + }) + .collect() +} pub fn load() -> CatalogAuth { - CatalogAuth { - access_key: read_storage(ACCESS_KEY_STORAGE).unwrap_or_default(), - secret_key: read_storage(SECRET_KEY_STORAGE).unwrap_or_default(), + let values = accounts(); + read_storage(ACTIVE_STORAGE) + .and_then(|v| v.parse::().ok()) + .and_then(|i| values.get(i).cloned()) + .unwrap_or_else(|| values.into_iter().next().unwrap_or_default()) +} +pub fn select(index: usize) { + write_storage(ACTIVE_STORAGE, &index.to_string()); +} +pub fn save_at(selected: Option, auth: &CatalogAuth) -> usize { + let auth = CatalogAuth { + label: auth.label.trim().to_owned(), + access_key: auth.access_key.trim().to_owned(), + secret_key: auth.secret_key.trim().to_owned(), + }; + let mut values = accounts(); + let index = selected.filter(|i| *i < values.len()).unwrap_or_else(|| { + values + .iter() + .position(|item| item.access_key == auth.access_key) + .unwrap_or(values.len()) + }); + if index == values.len() { + values.push(auth.clone()); + } else { + values[index] = auth.clone(); + } + select(index); + if let Ok(json) = serde_json::to_string(&values) { + write_storage(ACCOUNTS_STORAGE, &json); } + write_storage(ACCESS_KEY_STORAGE, &auth.access_key); + write_storage(SECRET_KEY_STORAGE, &auth.secret_key); + index } - -pub fn save(auth: &CatalogAuth) { - write_storage(ACCESS_KEY_STORAGE, auth.access_key.trim()); - write_storage(SECRET_KEY_STORAGE, auth.secret_key.trim()); +pub fn remove(index: usize) { + let mut values = accounts(); + if index >= values.len() { + return; + } + values.remove(index); + if let Ok(json) = serde_json::to_string(&values) { + write_storage(ACCOUNTS_STORAGE, &json); + } + if values.is_empty() { + write_storage(ACTIVE_STORAGE, "0"); + write_storage(ACCESS_KEY_STORAGE, ""); + write_storage(SECRET_KEY_STORAGE, ""); + } else { + select(index.min(values.len() - 1)); + } } - pub fn credentials() -> Option<(String, String)> { let auth = load(); auth.is_configured().then(|| { @@ -36,41 +102,37 @@ pub fn credentials() -> Option<(String, String)> { ) }) } - fn read_storage(key: &str) -> Option { let window = web_sys::window()?; let storage = window.local_storage().ok().flatten()?; storage.get_item(key).ok().flatten() } - fn write_storage(key: &str, value: &str) { - let Some(window) = web_sys::window() else { - return; - }; - let Some(storage) = window.local_storage().ok().flatten() else { - return; - }; - let _ = storage.set_item(key, value); + if let Some(window) = web_sys::window() + && let Some(storage) = window.local_storage().ok().flatten() + { + let _ = storage.set_item(key, value); + } } - #[cfg(test)] mod tests { use super::*; - #[test] fn configured_requires_both_keys() { assert!(!CatalogAuth::default().is_configured()); assert!( !CatalogAuth { + label: String::new(), access_key: "ak".into(), - secret_key: String::new(), + secret_key: String::new() } .is_configured() ); assert!( CatalogAuth { + label: String::new(), access_key: "ak".into(), - secret_key: "sk".into(), + secret_key: "sk".into() } .is_configured() ); diff --git a/pchronicle-web/src/llm_settings.rs b/pchronicle-web/src/llm_settings.rs index 298d389d8..e1845115e 100644 --- a/pchronicle-web/src/llm_settings.rs +++ b/pchronicle-web/src/llm_settings.rs @@ -13,8 +13,24 @@ pub fn LlmSettings( let mut api_key = use_signal(|| config.api_key.clone()); let mut model = use_signal(|| config.model.clone()); let initial_catalog = catalog_auth::load(); + let saved_accounts = catalog_auth::accounts(); + let mut catalog_accounts = use_signal(|| saved_accounts); + let mut editing = use_signal(|| { + let accounts = catalog_auth::accounts(); + (!accounts.is_empty()).then(|| { + accounts + .iter() + .position(|a| a == &initial_catalog) + .unwrap_or(0) + }) + }); + let mut catalog_label = use_signal(|| initial_catalog.label.clone()); let mut catalog_access_key = use_signal(|| initial_catalog.access_key.clone()); let mut catalog_secret_key = use_signal(|| initial_catalog.secret_key.clone()); + let accounts = catalog_accounts(); + let selected_profile = editing() + .map(|i| i.to_string()) + .unwrap_or_else(|| "__new__".into()); rsx! { div { class: "pc2-modal-backdrop high", section { class: "pc2-settings", role: "dialog", aria_modal: "true", @@ -26,7 +42,18 @@ pub fn LlmSettings( "Catalog keys are sent to this pChronicle server as request headers so it can authorize queries. Assistant keys stay in this browser and are sent only to the OpenAI-compatible endpoint." } div { class: "pc2-form", - p { class: "eyebrow", "Catalog" } + p { class: "eyebrow", "Catalog identity" } + label { span { "Profile" } + select { value: "{selected_profile}", onchange: move |event| { + let value = event.value(); + if value == "__new__" { editing.set(None); catalog_label.set(String::new()); catalog_access_key.set(String::new()); catalog_secret_key.set(String::new()); } + else if let Ok(index) = value.parse::() { catalog_auth::select(index); let values = catalog_auth::accounts(); if let Some(auth) = values.get(index) { editing.set(Some(index)); catalog_label.set(auth.label.clone()); catalog_access_key.set(auth.access_key.clone()); catalog_secret_key.set(auth.secret_key.clone()); } } + }, + option { value: "__new__", "+ New profile" } + for (index, auth) in accounts.iter().enumerate() { option { value: "{index}", "{auth.label}" } } + } + } + label { span { "Profile name" } input { placeholder: "e.g. Production", value: "{catalog_label}", oninput: move |event| catalog_label.set(event.value()) } } label { span { "Access key" } input { value: "{catalog_access_key}", oninput: move |event| catalog_access_key.set(event.value()) } } label { span { "Secret key" } input { r#type: "password", value: "{catalog_secret_key}", oninput: move |event| catalog_secret_key.set(event.value()) } } p { class: "eyebrow", "Assistant model" } @@ -36,13 +63,19 @@ pub fn LlmSettings( } footer { button { class: "button", onclick: on_close, "Cancel" } + if editing().is_some() { + button { class: "button danger", onclick: move |_| { if let Some(index) = editing() { catalog_auth::remove(index); let values = catalog_auth::accounts(); catalog_accounts.set(values.clone()); if let Some(auth) = values.first() { editing.set(Some(0)); catalog_label.set(auth.label.clone()); catalog_access_key.set(auth.access_key.clone()); catalog_secret_key.set(auth.secret_key.clone()); } else { editing.set(None); catalog_label.set(String::new()); catalog_access_key.set(String::new()); catalog_secret_key.set(String::new()); } } }, "Delete profile" } + } button { class: "button primary", onclick: move |_| { - catalog_auth::save(&CatalogAuth { + let index = catalog_auth::save_at(editing(), &CatalogAuth { + label: catalog_label(), access_key: catalog_access_key(), secret_key: catalog_secret_key(), }); + catalog_accounts.set(catalog_auth::accounts()); + editing.set(Some(index)); on_save.call(LlmConfig { api_base: api_base(), api_key: api_key(), diff --git a/pchronicle-web/src/model.rs b/pchronicle-web/src/model.rs index e1c899f75..55fa79f34 100644 --- a/pchronicle-web/src/model.rs +++ b/pchronicle-web/src/model.rs @@ -135,6 +135,10 @@ pub struct CatalogTree { #[serde(default)] pub failed_count: usize, #[serde(default)] + pub dataset_count: Option, + #[serde(default)] + pub trajectory_count: Option, + #[serde(default)] pub ready_sources: Option, #[serde(default)] pub error_sources: Option, @@ -159,6 +163,10 @@ pub struct CatalogTreeChild { #[serde(default)] pub failed_count: usize, #[serde(default)] + pub dataset_count: Option, + #[serde(default)] + pub trajectory_count: Option, + #[serde(default)] pub total_tokens: Option, #[serde(default)] pub entries: Vec, diff --git a/pchronicle-web/src/workspace.rs b/pchronicle-web/src/workspace.rs index cabf2195c..9598a9148 100644 --- a/pchronicle-web/src/workspace.rs +++ b/pchronicle-web/src/workspace.rs @@ -10,6 +10,7 @@ use wasm_bindgen::closure::Closure; use crate::agent::{self, ThreadMessage, ThreadRole}; use crate::api; use crate::catalog::CatalogExplorer; +use crate::catalog_auth; use crate::chat_view::normalize_trace_view; use crate::components::{ DataTable, HighlightedText, RichBlock, StepDrawer, TrajectoryView, WorkspaceIcon, @@ -283,6 +284,7 @@ pub fn App() -> Element { let catalog_loading = use_signal(|| false); let mut offset = use_signal(|| 0usize); let mut error = use_signal(|| None::); + let mut catalog_auth_configured = use_signal(|| catalog_auth::load().is_configured()); let mut selected_run = use_signal(move || initial_run); let mut analysis = use_signal(|| None::); @@ -543,7 +545,21 @@ pub fn App() -> Element { button { class: if copilot_open() { "rail-button active" } else { "rail-button" }, aria_label: "Toggle Assistant", aria_expanded: copilot_open(), onclick: move |_| copilot_open.set(!copilot_open()), WorkspaceIcon { name: "assistant" } span { {ASSISTANT} } } button { class: if settings_open() { "rail-button active" } else { "rail-button" }, aria_label: "Settings", onclick: move |_| settings_open.set(true), WorkspaceIcon { name: "keys" } span { "Keys" } } } - div { class: "rail-status", span { class: "live-dot" } span { "Local" } } + { + let identity = catalog_auth::load(); + let configured = identity.is_configured(); + let label = if configured { + if identity.label.is_empty() { "Catalog".to_string() } else { identity.label.clone() } + } else { + "Public access".to_string() + }; + rsx! { + div { class: "rail-status", title: "Active catalog profile: {label}", + span { class: if configured { "live-dot" } else { "live-dot muted" } } + span { "Local · {label}" } + } + } + } } main { id: "pc2-main", class: "pc2-main", tabindex: "-1", @@ -558,6 +574,8 @@ pub fn App() -> Element { CatalogExplorer { tree: catalog_tree(), loading: catalog_loading(), + auth_required: !catalog_auth_configured() && catalog_tree().is_none(), + on_settings: move |_| settings_open.set(true), on_open: move |(dataset, prefix): (String, String)| { catalog_dataset.set(dataset); catalog_prefix.set(prefix); @@ -893,6 +911,7 @@ pub fn App() -> Element { llm::save_config(&value); llm_config.set(value); settings_open.set(false); + catalog_auth_configured.set(catalog_auth::load().is_configured()); }, } } @@ -1071,6 +1090,16 @@ fn load_catalog_tree( spawn(async move { match api::explorer_tree(&dataset, &prefix).await { Ok(value) => tree.set(Some(value)), + Err(failure) + if matches!(failure.status, 400 | 401) + && dataset.is_empty() + && prefix.is_empty() => + { + match api::explorer_tree_anonymous(&dataset, &prefix).await { + Ok(value) => tree.set(Some(value)), + Err(failure) => error.set(Some(workspace_notice(&failure))), + } + } Err(failure) => error.set(Some(workspace_notice(&failure))), } loading.set(false);