diff --git a/crates/persisting-pchronicle-cli/src/server/acceleration.rs b/crates/persisting-pchronicle-cli/src/server/acceleration.rs index a1ef4b1dd..498fc098a 100644 --- a/crates/persisting-pchronicle-cli/src/server/acceleration.rs +++ b/crates/persisting-pchronicle-cli/src/server/acceleration.rs @@ -1078,6 +1078,16 @@ async fn bounded_summary_jsonl( max_rows: u64, max_bytes: usize, ) -> Result { + if tracing::enabled!(target: "pchronicle.query", tracing::Level::DEBUG) { + match engine.query_jsonl(&format!("EXPLAIN {sql}")).await { + Ok(plan) => { + tracing::debug!(target: "pchronicle.query", sql = %sql, plan = %plan, "summary query plan") + } + Err(error) => { + tracing::debug!(target: "pchronicle.query", sql = %sql, error = %error, "summary query plan unavailable") + } + } + } let mut output = super::BoundedOutput::new(max_bytes); let result = engine .write_query_jsonl_with_max_rows(sql, &mut output, Some(max_rows)) @@ -1087,6 +1097,13 @@ async fn bounded_summary_jsonl( "run summary byte budget exhausted; narrow the dataset or source scope" ); result.context("stream run summary within its row budget")?; + tracing::debug!( + target: "pchronicle.query", + sql = %sql, + rows = output.bytes.iter().filter(|byte| **byte == b'\n').count(), + bytes = output.bytes.len(), + "summary query completed" + ); String::from_utf8(output.bytes).context("decode run summary JSONL") } diff --git a/crates/persisting-pchronicle-cli/src/server/mod.rs b/crates/persisting-pchronicle-cli/src/server/mod.rs index 37beb78d7..3ce0dc794 100644 --- a/crates/persisting-pchronicle-cli/src/server/mod.rs +++ b/crates/persisting-pchronicle-cli/src/server/mod.rs @@ -42,7 +42,7 @@ use acceleration::{AccelerationStatus, ServerAcceleration}; use problem::{ ApiError, CHAIN_LIMIT, LOG_TARGET, QUERY_LOG_LIMIT, ROOT_CAUSE_LIMIT, truncate_utf8, }; -use request_log::{FtsDiagnostics, RequestId}; +use request_log::{FtsDiagnostics, RequestId, RequestMetrics}; fn fail(request_id: &RequestId, handler: &'static str, error: anyhow::Error) -> ApiError { ApiError::from_anyhow(request_id.as_str(), handler, error) @@ -55,7 +55,8 @@ use problem::BoundaryCode; struct AppState { config: Arc, catalog: Arc>>>, - catalog_refresh: Arc>, + /// Serializes global refreshes and stores the next automatic retry time. + catalog_refresh: Arc>, catalog_refresh_interval: Duration, trajectory_cache: Arc>>, /// Gateway-backed Warehouses read canonical events from the latest @@ -67,7 +68,7 @@ struct AppState { scoped_queries: Arc, } -const DEFAULT_CATALOG_REFRESH_INTERVAL: Duration = Duration::from_secs(5); +const DEFAULT_CATALOG_REFRESH_INTERVAL: Duration = Duration::from_secs(30); #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct HomeLink { @@ -209,7 +210,7 @@ fn app_state_with_catalog_refresh_interval( AppState { config: Arc::new(config), catalog: Arc::new(tokio::sync::RwLock::new(None)), - catalog_refresh: Arc::new(tokio::sync::Mutex::new(())), + catalog_refresh: Arc::new(tokio::sync::Mutex::new(Instant::now())), catalog_refresh_interval, trajectory_cache: Arc::new(tokio::sync::RwLock::new(None)), live_reads: false, @@ -301,10 +302,12 @@ impl PreparedWarehouse { let snapshot_id = runtime.snapshot.snapshot_id().to_string(); *self.state.catalog.write().await = Some(runtime); *self.state.trajectory_cache.write().await = None; + self.state.scoped_queries.invalidate(); snapshot_id } async fn refresh_runtime(&self) -> anyhow::Result> { + let _refresh = self.state.catalog_refresh.lock().await; let runtime = build_catalog_runtime(&self.state.config).await?; self.install_catalog_runtime(runtime.clone()).await; Ok(runtime) @@ -497,16 +500,35 @@ async fn build_catalog_runtime( async fn build_scoped_query_runtime( config: &ChronicleServerConfig, scope: persisting_pchronicle::storage::QueryScope, + cached_files: Vec, ) -> anyhow::Result> { - let snapshot = Arc::new( - DatasetCatalogSnapshot::discover_scoped( - config.datasets.clone(), - config.default_dataset.clone(), - config.catalog_options, - scope, - ) - .await?, + let target = match scope.source_file.clone() { + Some(file) => persisting_pchronicle::storage::ResolveTarget::Dataset { + mount: scope.dataset.clone(), + file, + }, + None => persisting_pchronicle::storage::ResolveTarget::Mount { + mount: scope.dataset.clone(), + prefix: None, + }, + }; + let cached_datasets: Vec<_> = cached_files + .iter() + .map(|file| { + persisting_pchronicle::storage::CachedDataset::new(scope.dataset.clone(), file.clone()) + }) + .collect(); + let mode = if cached_datasets.is_empty() { + persisting_pchronicle::storage::ResolveMode::Fresh + } else { + persisting_pchronicle::storage::ResolveMode::Cached + }; + let resolver = persisting_pchronicle::storage::DatasetResolver::new( + config.datasets.clone(), + config.default_dataset.clone(), + config.catalog_options, ); + let snapshot = Arc::new(resolver.resolve(target, mode, &cached_datasets).await?); let engine = Arc::new( snapshot .clone() @@ -547,16 +569,50 @@ async fn current_catalog_for_runs( if runtime.built_at.elapsed() < state.catalog_refresh_interval { return Ok(runtime); } - let _refresh = state.catalog_refresh.lock().await; - let current = state.catalog.read().await.clone().unwrap_or(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 { + let _refresh = state.catalog_refresh.lock().await; + return Ok(rebuild_catalog_for_runs(state, runtime).await); + } + if let Ok(mut refresh) = state.catalog_refresh.clone().try_lock_owned() + && *refresh <= Instant::now() + && let Ok(permit) = query_admission::REFRESH_SLOT.try_acquire() + { + let state = state.clone(); + let current = runtime.clone(); + tokio::spawn(async move { + let _permit = permit; + rebuild_catalog_for_runs(&state, current).await; + *refresh = Instant::now() + state.catalog_refresh_interval; + }); + } + Ok(runtime) +} + +async fn rebuild_catalog_for_runs( + state: &AppState, + current: Arc, +) -> Arc { + let current = state.catalog.read().await.clone().unwrap_or(current); if current.built_at.elapsed() < state.catalog_refresh_interval { - return Ok(current); + return current; + } + // Publish only after summaries are readable too: discovery alone can succeed + // while a lazily opened source is broken. + let refreshed = async { + let runtime = build_catalog_runtime(&state.config).await?; + runtime + .acceleration + .run_summaries(&runtime.snapshot, &runtime.engine) + .await?; + anyhow::Ok(runtime) } - match build_catalog_runtime(&state.config).await { + .await; + match refreshed { Ok(runtime) => { *state.catalog.write().await = Some(Arc::clone(&runtime)); *state.trajectory_cache.write().await = None; - Ok(runtime) + runtime } Err(error) => { tracing::warn!( @@ -565,7 +621,7 @@ async fn current_catalog_for_runs( chain = %truncate_utf8(&format!("{error:#}"), CHAIN_LIMIT), "automatic Catalog refresh failed; retaining the last valid snapshot" ); - Ok(current) + current } } } @@ -596,8 +652,11 @@ fn catalog_response(state: &AppState, runtime: &CatalogRuntime) -> CatalogRespon async fn catalog( State(state): State, request_id: RequestId, + metrics: RequestMetrics, ) -> Result, ApiError> { + let started = Instant::now(); let runtime = current_catalog(&state, &request_id).await?; + metrics.record("catalog", started); Ok(Json(catalog_response(&state, &runtime))) } @@ -618,10 +677,12 @@ async fn refresh_catalog( async fn runs( State(state): State, request_id: RequestId, + metrics: RequestMetrics, ) -> Result>, ApiError> { - Ok(Json( - load_run_summaries(&state, None, None, &request_id).await?, - )) + let started = Instant::now(); + let summaries = load_run_summaries(&state, None, None, &request_id, Some(&metrics)).await?; + metrics.record("summary_total", started); + Ok(Json(summaries)) } async fn load_run_summaries( @@ -629,6 +690,7 @@ async fn load_run_summaries( dataset: Option<&str>, file: Option<&str>, request_id: &RequestId, + metrics: Option<&RequestMetrics>, ) -> Result, ApiError> { let file = file.map(str::trim).filter(|value| !value.is_empty()); if let Some(dataset) = dataset { @@ -636,24 +698,84 @@ async fn load_run_summaries( dataset: dataset.to_owned(), source_file: file.map(str::to_owned), }; - return state - .scoped_queries - .run(scope.clone(), || async { - let runtime = build_scoped_query_runtime(&state.config, scope).await?; - runtime - .acceleration - .scoped_run_summaries(&runtime.snapshot, &runtime.engine, Some(dataset), file) - .await + let config = state.config.clone(); + let query_scope = scope.clone(); + let cached_files = if file.is_none() && !state.live_reads && !state.catalog_query_worker { + browse_coordinator(state) + .await + .cached_source_paths( + config + .datasets + .iter() + .find(|mount| mount.name == dataset) + .expect("validated dataset mount"), + ) + .await + } else { + Vec::new() + }; + let query_metrics = metrics.cloned(); + let execute = move |background: bool| async move { + // Background work must not be attributed to the triggering HTTP request. + let metrics = query_metrics.filter(|_| !background); + let phase = Instant::now(); + let runtime = + build_scoped_query_runtime(&config, query_scope.clone(), cached_files.clone()) + .await?; + if let Some(metrics) = &metrics { + metrics.record("summary_catalog", phase); + } + let phase = Instant::now(); + let summaries = runtime + .acceleration + .scoped_run_summaries( + &runtime.snapshot, + &runtime.engine, + Some(&query_scope.dataset), + query_scope.source_file.as_deref(), + ) + .await; + if let Some(metrics) = &metrics { + metrics.record("summary_sql", phase); + } + summaries + }; + let started = Instant::now(); + let result = if state.live_reads || state.catalog_query_worker { + state + .scoped_queries + .run(scope, || execute(false)) + .await + .map(|summaries| (summaries, "summary_uncached")) + } else { + state.scoped_queries.cached(scope, execute).await + }; + let result = result + .map(|(summaries, cache_status)| { + if let Some(metrics) = metrics { + metrics.record(cache_status, started); + } + summaries.as_ref().clone() }) - .await - .map(|summaries| summaries.as_ref().clone()) .map_err(|error| fail(request_id, "load_run_summaries", error)); + if let Some(metrics) = metrics { + metrics.record("summary_admission", started); + } + return result; } + let phase = Instant::now(); let runtime = current_catalog_for_runs(state, request_id).await?; + if let Some(metrics) = metrics { + metrics.record("summary_catalog", phase); + } + let phase = Instant::now(); let summaries = runtime .acceleration .scoped_run_summaries(&runtime.snapshot, &runtime.engine, dataset, file) .await; + if let Some(metrics) = metrics { + metrics.record("summary_sql", phase); + } summaries .map(|summaries| summaries.as_ref().clone()) .map_err(|error| fail(request_id, "load_run_summaries", error)) @@ -986,6 +1108,7 @@ async fn try_on_demand_storyline_runs_page( async fn explorer_runs( State(state): State, request_id: RequestId, + metrics: RequestMetrics, fts: FtsDiagnostics, query: Result, QueryRejection>, ) -> Result, ApiError> { @@ -1001,8 +1124,16 @@ async fn explorer_runs( .as_deref() .map(str::trim) .filter(|value| !value.is_empty() && *value != "all"); - let summaries = - load_run_summaries(&state, dataset_filter, query.file.as_deref(), &request_id).await?; + let started = Instant::now(); + let summaries = load_run_summaries( + &state, + dataset_filter, + query.file.as_deref(), + &request_id, + Some(&metrics), + ) + .await?; + metrics.record("summary_total", started); let (fts_matches, fts_available, search_mode) = if query .q .as_deref() @@ -1234,6 +1365,7 @@ fn preview_needle(query: &str) -> String { async fn explorer_tree( State(state): State, request_id: RequestId, + metrics: RequestMetrics, query: Result, QueryRejection>, ) -> Result, ApiError> { let query = api_query(query)?; @@ -1244,10 +1376,12 @@ async fn explorer_tree( .filter(|value| !value.is_empty()); let prefix = query.prefix.as_deref().unwrap_or(""); let Some(name) = dataset else { + let started = Instant::now(); let view = browse_coordinator(&state) .await .roots(&state.config.datasets) .await; + metrics.record("browse", started); return Ok(Json(serde_json::to_value(view).unwrap())); }; let Some(mount) = state @@ -1260,11 +1394,13 @@ async fn explorer_tree( serde_json::to_value(explorer::catalog_tree_from_path_list(name, prefix, &[])).unwrap(), )); }; + let started = Instant::now(); let view = browse_coordinator(&state) .await .tree(mount, prefix) .await .map_err(|error| fail(&request_id, "explorer_tree", error))?; + metrics.record("browse", started); Ok(Json(serde_json::to_value(view).unwrap())) } @@ -1272,45 +1408,71 @@ async fn resolve_run_summary( state: &AppState, query: &SessionQuery, request_id: &RequestId, + metrics: Option<&RequestMetrics>, ) -> Result { + // Exact Storyline leaves can identify a run with one document lookup; + // avoid scanning the whole source just to build a summary row. + if query + .dataset + .as_deref() + .is_some_and(|value| !value.trim().is_empty() && value != "all") + && query + .file + .as_deref() + .is_some_and(|value| !value.trim().is_empty()) + && !query.session_id.trim().is_empty() + && let Some(run) = try_resolve_on_demand_storyline_run(state, query, request_id).await? + { + return Ok(run); + } let dataset_filter = query .dataset .as_deref() .map(str::trim) .filter(|value| !value.is_empty() && *value != "all"); - let mut matches = load_run_summaries(state, dataset_filter, query.file.as_deref(), request_id) - .await? - .into_iter() - .filter(|run| { - query - .dataset + let mut matches = load_run_summaries( + state, + dataset_filter, + query.file.as_deref(), + request_id, + metrics, + ) + .await? + .into_iter() + .filter(|run| { + query + .dataset + .as_ref() + .filter(|value| !value.is_empty()) + .is_none_or(|value| value == &run.dataset) + && query + .file .as_ref() .filter(|value| !value.is_empty()) - .is_none_or(|value| value == &run.dataset) - && query - .file - .as_ref() - .filter(|value| !value.is_empty()) - .is_none_or(|value| value == &run.file) - && query - .run_id - .as_ref() - .filter(|value| !value.is_empty()) - .is_none_or(|value| run.run_id.as_ref() == Some(value)) - && (run.agent_id == query.agent_id - || run.model_name.as_deref() == Some(query.agent_id.as_str())) - && run.session_id == query.session_id - }) - .collect::>(); + .is_none_or(|value| value == &run.file) + && query + .run_id + .as_ref() + .filter(|value| !value.is_empty()) + .is_none_or(|value| run.run_id.as_ref() == Some(value)) + && (run.agent_id == query.agent_id + || run.model_name.as_deref() == Some(query.agent_id.as_str())) + && run.session_id == query.session_id + }) + .collect::>(); // Legacy browser URLs did not carry the Catalog key. Treat the old // root_session_id coordinate as an ambiguity breaker, not as a required // identity field: direct JSON sources may not preserve that field in the // normalized run row. + let filter_started = Instant::now(); if matches.len() > 1 && let Some(root) = &query.root_session_id { matches.retain(|run| run.root_session_id.as_ref() == Some(root)); } + if let Some(metrics) = metrics { + metrics.record("resolve_filter", filter_started); + } if matches.is_empty() { if let Some(run) = try_resolve_on_demand_storyline_run(state, query, request_id).await? { return Ok(run); @@ -1357,9 +1519,13 @@ async fn try_resolve_on_demand_storyline_run( let Some(dataset) = runtime.snapshot.dataset(dataset_name) else { return Ok(None); }; + let exact_source = dataset.sources.iter().any(|source| { + source.kind != persisting_pchronicle::storage::CatalogSourceKind::Directory + && source.file == file + }); if dataset.sources.iter().any(|source| { source.kind != persisting_pchronicle::storage::CatalogSourceKind::Directory - && (source.file == file || source.file.starts_with(&format!("{file}/"))) + && source.file.starts_with(&format!("{file}/")) }) { return Ok(None); } @@ -1367,7 +1533,7 @@ async fn try_resolve_on_demand_storyline_run( source.kind == persisting_pchronicle::storage::CatalogSourceKind::Directory && (file == source.file || file.starts_with(&format!("{}/", source.file))) }); - if !under_directory && !dataset.sources.is_empty() { + if !under_directory && !exact_source && !dataset.sources.is_empty() { return Ok(None); } @@ -1553,7 +1719,7 @@ async fn load_events( query: &SessionQuery, request_id: &RequestId, ) -> Result { - let run = resolve_run_summary(state, query, request_id).await?; + let run = resolve_run_summary(state, query, request_id, None).await?; let bundle = catalog_or_on_demand_trajectory_bundle(state, &run, request_id, "load_events").await?; let document = bundle.event_view; @@ -1616,7 +1782,7 @@ async fn storyline( query: Result, QueryRejection>, ) -> Result, ApiError> { let query = api_query(query)?; - let run = resolve_run_summary(&state, &query, &request_id).await?; + let run = resolve_run_summary(&state, &query, &request_id, None).await?; let bundle = catalog_or_on_demand_trajectory_bundle(&state, &run, &request_id, "storyline").await?; Ok(Json( @@ -1755,9 +1921,14 @@ async fn load_trajectory( state: &AppState, query: &SessionQuery, request_id: &RequestId, + metrics: &RequestMetrics, ) -> Result { - let run = resolve_run_summary(state, query, request_id).await?; + let phase = Instant::now(); + let run = resolve_run_summary(state, query, request_id, Some(metrics)).await?; + metrics.record("resolve", phase); + let phase = Instant::now(); let runtime = current_catalog(state, request_id).await?; + metrics.record("catalog", phase); let cache_key = format!( "{}\u{1f}{}\u{1f}{}\u{1f}{}", runtime.snapshot.snapshot_id(), @@ -1773,6 +1944,7 @@ async fn load_trajectory( .as_ref() .filter(|(key, _)| key == &cache_key) { + metrics.record("trajectory_cache", Instant::now()); return Ok(loaded.clone()); } if run.format.as_deref() == Some("compact-jsonl/v1") { @@ -1783,8 +1955,10 @@ async fn load_trajectory( turns: Vec::new(), }); } + let phase = Instant::now(); let bundle = catalog_or_on_demand_trajectory_bundle(state, &run, request_id, "load_trajectory").await?; + metrics.record("trajectory_read", phase); let event_provenance = bundle.event_view.provenance; let records = bundle.event_view.document.events; let document = bundle.storyline; @@ -1863,10 +2037,11 @@ async fn load_trajectory( async fn trajectory_view( State(state): State, request_id: RequestId, + metrics: RequestMetrics, query: Result, QueryRejection>, ) -> Result, ApiError> { let query = api_query(query)?; - let loaded = load_trajectory(&state, &query, &request_id).await?; + let loaded = load_trajectory(&state, &query, &request_id, &metrics).await?; let mut event_kind_counts = BTreeMap::new(); for event in &loaded.records { *event_kind_counts.entry(event.kind.clone()).or_insert(0) += 1; @@ -1889,10 +2064,11 @@ async fn trajectory_view( async fn explorer_run( State(state): State, request_id: RequestId, + metrics: RequestMetrics, query: Result, QueryRejection>, ) -> Result, ApiError> { let query = api_query(query)?; - let loaded = load_trajectory(&state, &query, &request_id).await?; + let loaded = load_trajectory(&state, &query, &request_id, &metrics).await?; Ok(Json(explorer::analyze( loaded.run, &loaded.turns, @@ -1913,7 +2089,7 @@ async fn explorer_record( query: Result, QueryRejection>, ) -> Result, ApiError> { let query = api_query(query)?; - let run = resolve_run_summary(&state, &query, &request_id).await?; + let run = resolve_run_summary(&state, &query, &request_id, None).await?; if run.format.as_deref() != Some("compact-jsonl/v1") { return Err(ApiError::not_found("run is not a compact JSONL record")); } @@ -1960,13 +2136,16 @@ impl TurnsQuery { async fn explorer_turns( State(state): State, request_id: RequestId, + metrics: RequestMetrics, fts: FtsDiagnostics, query: Result, QueryRejection>, ) -> Result, ApiError> { let query = api_query(query)?; let session = query.session(); - let loaded = load_trajectory(&state, &session, &request_id).await?; + let loaded = load_trajectory(&state, &session, &request_id, &metrics).await?; + let phase = Instant::now(); let runtime = current_catalog(&state, &request_id).await?; + metrics.record("turn_catalog", phase); // Nested Directory Storylines are opened on-demand and are absent from the // prepared catalog; skip FTS path probing and keep in-memory turn pages. let paths = runtime @@ -1974,6 +2153,7 @@ async fn explorer_turns( .storyline_table_paths(&loaded.run.dataset, &loaded.run.file) .ok() .flatten(); + let phase = Instant::now(); let mut fts_available = if let Some(paths) = paths.as_ref() { match storyline_steps_fts_available(paths).await { Ok(available) => available, @@ -1985,6 +2165,7 @@ async fn explorer_turns( } else { false }; + metrics.record("turn_fts_probe", phase); let mut search_mode = if query .q .as_deref() @@ -2010,6 +2191,7 @@ async fn explorer_turns( .map_err(|error| ApiError::invalid_request(error.to_string()))? .ok_or_else(|| ApiError::invalid_request("search query must not be empty"))?; let runtime = current_catalog(&state, &request_id).await?; + let phase = Instant::now(); let (predicate, available, fts_errors) = crate::find_expression_predicate_for_dataset( &runtime.snapshot, &expression, @@ -2019,6 +2201,7 @@ async fn explorer_turns( .await .map_err(|error| fail(&request_id, "explorer_turns", error))?; fts.extend(fts_errors); + metrics.record("turn_search", phase); fts_available = fts_available || available; let turns = if expression.has_text() || expression.has_step_json() { let predicate = predicate.ok_or_else(|| { @@ -2074,7 +2257,8 @@ async fn explorer_turns( } else { (loaded.turns.clone(), query.q.as_deref()) }; - Ok(Json(explorer::turn_page_with_search( + let phase = Instant::now(); + let page = explorer::turn_page_with_search( &turns, &loaded.records, search_query, @@ -2086,7 +2270,9 @@ async fn explorer_turns( mode: search_mode, tokenizer: fts_available.then_some("jieba"), }, - ))) + ); + metrics.record("turn_projection", phase); + Ok(Json(page)) } #[derive(Debug, Deserialize)] @@ -2103,6 +2289,7 @@ struct TurnDetailQuery { async fn explorer_turn( State(state): State, request_id: RequestId, + metrics: RequestMetrics, query: Result, QueryRejection>, ) -> Result, ApiError> { let query = api_query(query)?; @@ -2116,7 +2303,7 @@ async fn explorer_turn( offset: None, limit: None, }; - let loaded = load_trajectory(&state, &session, &request_id).await?; + let loaded = load_trajectory(&state, &session, &request_id, &metrics).await?; let item = loaded .turns .iter() @@ -2371,10 +2558,173 @@ fn format_query_fields() -> Vec { ] } +fn query_table_summaries() -> Vec { + vec![ + QueryTableSummary { + name: "sources", + description: "One row per discovered run data source", + kind: "table", + grain: "source", + fields: source_query_fields(), + }, + QueryTableSummary { + name: "runs", + description: "One row per run across the complete data path", + kind: "table", + grain: "run", + fields: run_query_fields(), + }, + QueryTableSummary { + name: "steps", + description: "Ordered user, agent, and system steps for every run", + kind: "table", + grain: "step", + fields: step_query_fields(), + }, + QueryTableSummary { + name: "tool_calls", + description: "Structured tool calls joined to their run and step", + kind: "table", + grain: "tool call", + fields: tool_call_query_fields(), + }, + QueryTableSummary { + name: "trajectories", + description: "One complete run with ordered step and tool-call arrays", + kind: "view", + grain: "complete run", + fields: trajectory_query_fields(), + }, + QueryTableSummary { + name: "events", + description: "Recorded events; empty for sources without recorded events", + kind: "table", + grain: "event", + fields: event_query_fields(), + }, + QueryTableSummary { + name: "atif", + description: "ATIF documents exposed as id plus JSONB data", + kind: "view", + grain: "ATIF document", + fields: format_query_fields(), + }, + QueryTableSummary { + name: "storyline", + description: "Storyline documents exposed as id plus JSONB data", + kind: "view", + grain: "Storyline document", + fields: format_query_fields(), + }, + QueryTableSummary { + name: "actf", + description: "ACTF documents exposed as id plus JSONB data", + kind: "view", + grain: "ACTF document", + fields: format_query_fields(), + }, + QueryTableSummary { + name: "openai_msg", + description: "OpenAI message documents exposed as id plus JSONB data", + kind: "view", + grain: "OpenAI message", + fields: format_query_fields(), + }, + QueryTableSummary { + name: "codex", + description: "Codex documents exposed as id plus JSONB data", + kind: "view", + grain: "Codex document", + fields: format_query_fields(), + }, + QueryTableSummary { + name: "markdown", + description: "Markdown/AgenticMD documents exposed as id plus JSONB data", + kind: "view", + grain: "Markdown document", + fields: format_query_fields(), + }, + QueryTableSummary { + name: "claude", + description: "Claude Code documents exposed as id plus JSONB data", + kind: "view", + grain: "Claude Code document", + fields: format_query_fields(), + }, + ] +} + +#[derive(Debug, Deserialize, Default)] +struct QueryTablesQuery { + /// Accept both `ui=true` (new clients) and `ui=1` (cached clients). + #[serde(default)] + ui: Option, +} + +async fn ui_query_catalog(state: &AppState) -> Result, ApiError> { + let browse = browse_coordinator(state).await; + let default_name = state + .config + .default_dataset + .clone() + .or_else(|| { + state + .config + .datasets + .first() + .map(|mount| mount.name.clone()) + }) + .unwrap_or_else(|| DEFAULT_DATASET_NAME.into()); + let storage_path = state + .config + .datasets + .iter() + .find(|mount| mount.name == default_name) + .map(|mount| mount.uri.clone()) + .unwrap_or_default(); + let mut datasets = Vec::with_capacity(state.config.datasets.len()); + for mount in &state.config.datasets { + let tree = browse.cached_dataset(mount).await; + let (ready_sources, error_sources) = tree + .map(|tree| { + tree.children.iter().fold((0, 0), |(ready, failed), child| { + ( + ready + usize::from(child.kind == "file"), + failed + usize::from(child.failed_count > 0), + ) + }) + }) + .unwrap_or_default(); + datasets.push(QueryDatasetSummary { + name: mount.name.clone(), + uri: mount.uri.clone(), + ready_sources, + error_sources, + }); + } + Ok(Json(QueryCatalog { + snapshot_id: String::new(), + read_only: true, + database: default_name, + storage_path, + path_column: "_file_", + datasets, + tables: query_table_summaries(), + })) +} + async fn query_tables( State(state): State, request_id: RequestId, + Query(query): Query, ) -> Result, ApiError> { + if query + .ui + .as_deref() + .is_some_and(|value| value == "true" || value == "1") + { + return ui_query_catalog(&state).await; + } let runtime = current_catalog(&state, &request_id).await?; let database = runtime .snapshot @@ -2410,99 +2760,7 @@ async fn query_tables( error_sources: dataset.error_source_count(), }) .collect(), - tables: vec![ - QueryTableSummary { - name: "sources", - description: "One row per discovered run data source", - kind: "table", - grain: "source", - fields: source_query_fields(), - }, - QueryTableSummary { - name: "runs", - description: "One row per run across the complete data path", - kind: "table", - grain: "run", - fields: run_query_fields(), - }, - QueryTableSummary { - name: "steps", - description: "Ordered user, agent, and system steps for every run", - kind: "table", - grain: "step", - fields: step_query_fields(), - }, - QueryTableSummary { - name: "tool_calls", - description: "Structured tool calls joined to their run and step", - kind: "table", - grain: "tool call", - fields: tool_call_query_fields(), - }, - QueryTableSummary { - name: "trajectories", - description: "One complete run with ordered step and tool-call arrays", - kind: "view", - grain: "complete run", - fields: trajectory_query_fields(), - }, - QueryTableSummary { - name: "events", - description: "Recorded events; empty for sources without recorded events", - kind: "table", - grain: "event", - fields: event_query_fields(), - }, - QueryTableSummary { - name: "atif", - description: "ATIF documents exposed as id plus JSONB data", - kind: "view", - grain: "ATIF document", - fields: format_query_fields(), - }, - QueryTableSummary { - name: "storyline", - description: "Storyline documents exposed as id plus JSONB data", - kind: "view", - grain: "Storyline document", - fields: format_query_fields(), - }, - QueryTableSummary { - name: "actf", - description: "ACTF documents exposed as id plus JSONB data", - kind: "view", - grain: "ACTF document", - fields: format_query_fields(), - }, - QueryTableSummary { - name: "openai_msg", - description: "OpenAI message documents exposed as id plus JSONB data", - kind: "view", - grain: "OpenAI message", - fields: format_query_fields(), - }, - QueryTableSummary { - name: "codex", - description: "Codex documents exposed as id plus JSONB data", - kind: "view", - grain: "Codex document", - fields: format_query_fields(), - }, - QueryTableSummary { - name: "markdown", - description: "Markdown/AgenticMD documents exposed as id plus JSONB data", - kind: "view", - grain: "Markdown document", - fields: format_query_fields(), - }, - QueryTableSummary { - name: "claude", - description: "Claude Code documents exposed as id plus JSONB data", - kind: "view", - grain: "Claude Code document", - fields: format_query_fields(), - }, - ], + tables: query_table_summaries(), })) } @@ -2666,7 +2924,7 @@ async fn compile_analysis( engine_detail: None, }, })?; - if request.snapshot_id != runtime.snapshot.snapshot_id() { + if request.snapshot_id != "ui-cache" && request.snapshot_id != runtime.snapshot.snapshot_id() { return Err(CompileHttpError::stale_snapshot()); } let schema = runtime diff --git a/crates/persisting-pchronicle-cli/src/server/query_admission.rs b/crates/persisting-pchronicle-cli/src/server/query_admission.rs index dc96c6270..3c2002a1a 100644 --- a/crates/persisting-pchronicle-cli/src/server/query_admission.rs +++ b/crates/persisting-pchronicle-cli/src/server/query_admission.rs @@ -1,9 +1,10 @@ -//! Coalesce overlapping scoped summary requests, without caching pinned sources -//! across independent requests. Admission covers discovery AND query execution. +//! Bound and coalesce UI summary queries. Cache successful summaries (not query +//! engines or source pins), returning stale results while one refresh runs. use std::collections::HashMap; use std::future::Future; use std::sync::{Arc, Mutex, Weak}; +use std::time::{Duration, Instant}; use anyhow::Result; use persisting_pchronicle::storage::QueryScope; @@ -13,6 +14,52 @@ use super::RunSummary; use super::acceleration::SharedAccelerationFailure; type Summaries = Arc>; +const REFRESH_INTERVAL: Duration = Duration::from_secs(30); +const MAX_CACHE_ENTRIES: usize = 32; +const MAX_CACHE_BYTES: usize = 64 * 1024 * 1024; +// One background summary refresh per process, with no waiting task queue. +pub(super) static REFRESH_SLOT: Semaphore = Semaphore::const_new(1); + +struct CachedSummaries { + summaries: Summaries, + built_at: Instant, + refresh_after: Instant, + last_used: Instant, + bytes: usize, +} + +#[derive(Default)] +struct SummaryCache { + generation: u64, + entries: HashMap, +} + +fn summary_bytes(summaries: &Vec) -> usize { + summaries.capacity() * std::mem::size_of::() + + summaries + .iter() + .map(|run| { + run.dataset.capacity() + + run.file.capacity() + + run.document_id.capacity() + + run.agent_id.capacity() + + run.session_id.capacity() + + run.path.capacity() + + run.status.capacity() + + [ + &run.run_id, + &run.model_name, + &run.root_session_id, + &run.format, + ] + .into_iter() + .flatten() + .map(String::capacity) + .sum::() + }) + .sum::() +} + #[derive(Default)] struct Flight { result: OnceCell>, @@ -21,6 +68,7 @@ struct Flight { pub(super) struct ScopedQueries { flights: Mutex>>, slots: Semaphore, + cache: Mutex, } impl Default for ScopedQueries { @@ -28,11 +76,119 @@ impl Default for ScopedQueries { Self { flights: Mutex::new(HashMap::new()), slots: Semaphore::new(2), + cache: Mutex::new(SummaryCache::default()), } } } impl ScopedQueries { + pub(super) fn invalidate(&self) { + let mut cache = self.cache.lock().unwrap(); + cache.generation += 1; + cache.entries.clear(); + // Requests started after explicit refresh must not join an older build. + self.flights.lock().unwrap().clear(); + } + + fn publish(&self, scope: QueryScope, generation: u64, summaries: Summaries) { + let bytes = summary_bytes(&summaries); + let mut cache = self.cache.lock().unwrap(); + if cache.generation != generation { + return; + } + cache.entries.remove(&scope); + if bytes > MAX_CACHE_BYTES { + return; + } + while cache.entries.len() >= MAX_CACHE_ENTRIES + || cache + .entries + .values() + .map(|entry| entry.bytes) + .sum::() + + bytes + > MAX_CACHE_BYTES + { + let oldest = cache + .entries + .iter() + .min_by_key(|(_, entry)| entry.last_used) + .map(|(key, _)| key.clone()) + .unwrap(); + cache.entries.remove(&oldest); + } + let now = Instant::now(); + cache.entries.insert( + scope, + CachedSummaries { + summaries, + built_at: now, + refresh_after: now + REFRESH_INTERVAL, + last_used: now, + bytes, + }, + ); + } + + pub(super) async fn cached( + self: &Arc, + scope: QueryScope, + execute: F, + ) -> Result<(Summaries, &'static str)> + where + F: FnOnce(bool) -> Fut + Send + 'static, + Fut: Future> + Send + 'static, + { + let (generation, cached) = { + let mut cache = self.cache.lock().unwrap(); + let generation = cache.generation; + let cached = cache.entries.get_mut(&scope).map(|entry| { + entry.last_used = Instant::now(); + let stale = entry.built_at.elapsed() >= REFRESH_INTERVAL; + let permit = (stale && entry.last_used >= entry.refresh_after) + .then(|| REFRESH_SLOT.try_acquire().ok()) + .flatten(); + if permit.is_some() { + entry.refresh_after = Instant::now() + REFRESH_INTERVAL; + } + (entry.summaries.clone(), stale, permit) + }); + (generation, cached) + }; + if let Some((summaries, stale, permit)) = cached { + if let Some(permit) = permit { + let queries = self.clone(); + tokio::spawn(async move { + let _permit = permit; + match queries.run(scope.clone(), || execute(true)).await { + Ok(summaries) => queries.publish(scope, generation, summaries), + Err(error) => { + let mut cache = queries.cache.lock().unwrap(); + if cache.generation == generation + && let Some(entry) = cache.entries.get_mut(&scope) + { + entry.refresh_after = Instant::now() + REFRESH_INTERVAL; + } + tracing::warn!(target: super::LOG_TARGET, dataset = %scope.dataset, + error = %error, "UI run summary refresh failed; retaining cached results"); + } + } + }); + } + return Ok(( + summaries, + if stale { + "summary_cache_stale" + } else { + "summary_cache_hit" + }, + )); + } + let summaries = self.run(scope.clone(), || execute(false)).await?; + self.publish(scope, generation, summaries.clone()); + Ok((summaries, "summary_cache_miss")) + } + pub(super) async fn run(&self, scope: QueryScope, execute: F) -> Result where F: FnOnce() -> Fut, @@ -77,6 +233,143 @@ mod tests { } } + fn expire(queries: &ScopedQueries, key: &QueryScope) { + let mut cache = queries.cache.lock().unwrap(); + let entry = cache.entries.get_mut(key).unwrap(); + entry.built_at = Instant::now() - REFRESH_INTERVAL; + entry.refresh_after = Instant::now(); + } + + #[tokio::test] + async fn ui_cache_reuses_results_bounds_refresh_and_retains_failures() { + let queries = Arc::new(ScopedQueries::default()); + let initial = Arc::new(Vec::new()); + let value = initial.clone(); + let (first, status) = queries + .cached(scope("a"), |_| async { Ok(value) }) + .await + .unwrap(); + assert_eq!(status, "summary_cache_miss"); + let (hit, status) = queries + .cached(scope("a"), |_| async { panic!("cache hit must not scan") }) + .await + .unwrap(); + assert!(Arc::ptr_eq(&first, &hit)); + assert_eq!(status, "summary_cache_hit"); + queries + .cached(scope("b"), |_| async { Ok(Arc::new(Vec::new())) }) + .await + .unwrap(); + expire(&queries, &scope("a")); + expire(&queries, &scope("b")); + let (started, ready) = tokio::sync::oneshot::channel(); + let (finish, wait) = tokio::sync::oneshot::channel(); + let updated = Arc::new(Vec::new()); + let value = updated.clone(); + let (stale, status) = queries + .cached(scope("a"), |_| async move { + started.send(()).unwrap(); + wait.await.unwrap(); + Ok(value) + }) + .await + .unwrap(); + assert!(Arc::ptr_eq(&initial, &stale)); + assert_eq!(status, "summary_cache_stale"); + ready.await.unwrap(); + for name in ["a", "b"] { + let (_, status) = queries + .cached(scope(name), |_| async { + panic!("only one background refresh") + }) + .await + .unwrap(); + assert_eq!(status, "summary_cache_stale"); + } + finish.send(()).unwrap(); + let permit = REFRESH_SLOT.acquire().await.unwrap(); + assert!(Arc::ptr_eq( + &queries.cache.lock().unwrap().entries[&scope("a")].summaries, + &updated + )); + drop(permit); + + expire(&queries, &scope("a")); + let (started, ready) = tokio::sync::oneshot::channel(); + queries + .cached(scope("a"), |_| async move { + started.send(()).unwrap(); + Err(anyhow::anyhow!("S3 unavailable")) + }) + .await + .unwrap(); + ready.await.unwrap(); + let permit = REFRESH_SLOT.acquire().await.unwrap(); + drop(permit); + let (retained, status) = queries + .cached(scope("a"), |_| async { panic!("retry must back off") }) + .await + .unwrap(); + assert!(Arc::ptr_eq(&retained, &updated)); + assert_eq!(status, "summary_cache_stale"); + + // Explicit refresh invalidates results and prevents a late background + // completion from republishing data from the previous generation. + expire(&queries, &scope("a")); + let (started, ready) = tokio::sync::oneshot::channel(); + let (finish, wait) = tokio::sync::oneshot::channel(); + queries + .cached(scope("a"), |_| async move { + started.send(()).unwrap(); + wait.await.unwrap(); + Ok(Arc::new(Vec::new())) + }) + .await + .unwrap(); + ready.await.unwrap(); + queries.invalidate(); + finish.send(()).unwrap(); + let _permit = REFRESH_SLOT.acquire().await.unwrap(); + assert!(queries.cache.lock().unwrap().entries.is_empty()); + } + + #[tokio::test] + async fn ui_cache_bounds_retention_and_never_caches_cold_failures() { + let queries = Arc::new(ScopedQueries::default()); + assert!( + queries + .cached(scope("a"), |_| async { + Err(anyhow::anyhow!("unavailable")) + }) + .await + .is_err() + ); + assert!(queries.cache.lock().unwrap().entries.is_empty()); + for n in 0..=MAX_CACHE_ENTRIES { + queries + .cached(scope(&n.to_string()), |_| async { + Ok(Arc::new(Vec::new())) + }) + .await + .unwrap(); + } + let cache = queries.cache.lock().unwrap(); + assert_eq!(cache.entries.len(), MAX_CACHE_ENTRIES); + assert!(!cache.entries.contains_key(&scope("0"))); + drop(cache); + // Account for allocated vector capacity, not just serialized row bytes. + let oversized = Vec::with_capacity(MAX_CACHE_BYTES / std::mem::size_of::() + 1); + queries.publish(scope("oversized"), 0, Arc::new(oversized)); + assert!( + !queries + .cache + .lock() + .unwrap() + .entries + .contains_key(&scope("oversized")) + ); + } + #[tokio::test] async fn overlapping_requests_share_work_but_later_request_repins() { let queries = ScopedQueries::default(); diff --git a/crates/persisting-pchronicle-cli/src/server/request_log.rs b/crates/persisting-pchronicle-cli/src/server/request_log.rs index 8a3487c14..13011a777 100644 --- a/crates/persisting-pchronicle-cli/src/server/request_log.rs +++ b/crates/persisting-pchronicle-cli/src/server/request_log.rs @@ -26,6 +26,27 @@ impl RequestId { #[derive(Clone, Default)] pub(crate) struct FtsDiagnostics(pub Arc>>); +#[derive(Clone, Default)] +pub(crate) struct RequestMetrics(pub Arc>>); + +impl RequestMetrics { + pub(crate) fn record(&self, name: &'static str, started: Instant) { + if let Ok(mut values) = self.0.lock() { + values.push((name, started.elapsed().as_millis() as u64)); + } + } + + fn server_timing(&self, total_ms: u64) -> String { + let mut values = self.0.lock().map(|v| v.clone()).unwrap_or_default(); + values.push(("total", total_ms)); + values + .into_iter() + .map(|(name, ms)| format!("{name};dur={ms}")) + .collect::>() + .join(", ") + } +} + impl FtsDiagnostics { pub(crate) fn push(&self, message: impl Into) { if let Ok(mut errors) = self.0.lock() { @@ -77,6 +98,21 @@ where } } +impl FromRequestParts for RequestMetrics +where + S: Send + Sync, +{ + type Rejection = std::convert::Infallible; + + async fn from_request_parts(parts: &mut Parts, _state: &S) -> Result { + Ok(parts + .extensions + .get::() + .cloned() + .unwrap_or_default()) + } +} + pub(crate) async fn warehouse_request_layer( mut request: Request, next: Next, @@ -92,10 +128,12 @@ pub(crate) async fn warehouse_request_layer( let query = request.uri().query().unwrap_or("").to_owned(); let started = Instant::now(); let fts = FtsDiagnostics::default(); + let metrics = RequestMetrics::default(); request .extensions_mut() .insert(RequestId(request_id.clone())); request.extensions_mut().insert(fts.clone()); + request.extensions_mut().insert(metrics.clone()); tracing::info!( target: LOG_TARGET, @@ -112,9 +150,10 @@ pub(crate) async fn warehouse_request_layer( .get::() .map(|value| value.0.clone()) .unwrap_or_default(); - let (response, error_fields) = attach_request_id(response, &request_id).await; - let elapsed_ms = u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX); + let (response, error_fields) = + attach_request_id(response, &request_id, &metrics, elapsed_ms).await; + let server_timing = metrics.server_timing(elapsed_ms); tracing::info!( target: LOG_TARGET, request_id = %request_id, @@ -122,6 +161,7 @@ pub(crate) async fn warehouse_request_layer( path = %path, status = status.as_u16(), elapsed_ms, + server_timing = %server_timing, query = %truncate_utf8(&query, QUERY_LOG_LIMIT), "warehouse request" ); @@ -135,6 +175,7 @@ pub(crate) async fn warehouse_request_layer( message = %message, root_cause = %root_cause, fts_errors = %fts_errors, + server_timing = %server_timing, "warehouse request rejected" ); } @@ -144,11 +185,16 @@ pub(crate) async fn warehouse_request_layer( async fn attach_request_id( response: Response, request_id: &str, + metrics: &RequestMetrics, + total_ms: u64, ) -> (Response, Option<(String, String)>) { let (mut parts, body) = response.into_parts(); 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 is_json = parts .headers .get(CONTENT_TYPE) diff --git a/crates/persisting-pchronicle-cli/src/server/tests.rs b/crates/persisting-pchronicle-cli/src/server/tests.rs index 02e21a05a..71c44a081 100644 --- a/crates/persisting-pchronicle-cli/src/server/tests.rs +++ b/crates/persisting-pchronicle-cli/src/server/tests.rs @@ -360,6 +360,13 @@ async fn middleware_echoes_request_id_on_json_errors() { response.headers().get("x-request-id").unwrap(), "client-id-123" ); + assert!( + response + .headers() + .get("server-timing") + .and_then(|value| value.to_str().ok()) + .is_some_and(|value| value.contains("total;dur=")) + ); let body = response_json(response).await; assert_eq!(body["request_id"], "client-id-123"); assert_eq!(body["code"], "invalid_request"); @@ -1570,6 +1577,130 @@ async fn warehouse_does_not_expose_unused_har_or_revisions_routes() { } } +#[tokio::test] +async fn unscoped_runs_refresh_in_background_and_back_off_after_failure() { + let root = tempfile::tempdir().unwrap(); + write_gateway_fixture(root.path(), "first.json", "first-session", "job"); + let config = ChronicleServerConfig::mounted(vec![ + DatasetMount::default(root.path().to_string_lossy().to_string()).unwrap(), + ]) + .unwrap(); + let state = app_state(config); + let request_id = RequestId("ui-cache-test".into()); + let mut initial = build_catalog_runtime(&state.config).await.unwrap(); + Arc::get_mut(&mut initial).unwrap().built_at = Instant::now() - Duration::from_secs(60); + let old_id = initial.snapshot.snapshot_id().to_owned(); + *state.catalog.write().await = Some(initial); + write_gateway_fixture(root.path(), "second.json", "second-session", "job"); + + // An ongoing Catalog refresh must never make a warm UI request wait. + let guard = state.catalog_refresh.lock().await; + let cached = tokio::time::timeout( + Duration::from_secs(1), + current_catalog_for_runs(&state, &request_id), + ) + .await + .unwrap() + .unwrap(); + assert_eq!(cached.snapshot.snapshot_id(), old_id); + drop(cached); + drop(guard); + let cached = current_catalog_for_runs(&state, &request_id).await.unwrap(); + assert_eq!(cached.snapshot.snapshot_id(), old_id); + drop(cached); + let permit = query_admission::REFRESH_SLOT.acquire().await.unwrap(); + let new_id = state + .catalog + .read() + .await + .as_ref() + .unwrap() + .snapshot + .snapshot_id() + .to_owned(); + assert_ne!(new_id, old_id); + drop(permit); + let summaries = load_run_summaries(&state, None, None, &request_id, None) + .await + .unwrap(); + assert_eq!(summaries.len(), 2); + + // A failed refresh must preserve the published runtime and delay retries. + std::fs::write(root.path().join("broken.json"), "{invalid").unwrap(); + { + let mut catalog = state.catalog.write().await; + Arc::get_mut(catalog.as_mut().unwrap()).unwrap().built_at = + Instant::now() - Duration::from_secs(60); + } + *state.catalog_refresh.lock().await = Instant::now(); + let retained = current_catalog_for_runs(&state, &request_id).await.unwrap(); + assert_eq!(retained.snapshot.snapshot_id(), new_id); + drop(retained); + let _permit = query_admission::REFRESH_SLOT.acquire().await.unwrap(); + assert_eq!( + state + .catalog + .read() + .await + .as_ref() + .unwrap() + .snapshot + .snapshot_id(), + new_id + ); + assert!(*state.catalog_refresh.lock().await > Instant::now()); +} + +#[tokio::test] +async fn explorer_runs_reuses_ui_summaries_until_explicit_catalog_refresh() { + use tower::ServiceExt; + + let root = tempfile::tempdir().unwrap(); + write_gateway_fixture(root.path(), "first.json", "first-session", "job"); + let app = router(root.path().to_string_lossy().to_string()); + let uri = "/api/explorer/runs?dataset=dataset&limit=1"; + let (status, initial) = get_json(&app, uri).await; + assert_eq!(status, StatusCode::OK, "{initial}"); + assert_eq!(initial["snapshot"]["total"], 1); + + write_gateway_fixture(root.path(), "second.json", "second-session", "job"); + let response = app + .clone() + .oneshot( + axum::http::Request::builder() + .uri(uri) + .body(axum::body::Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + assert!( + response.headers()["server-timing"] + .to_str() + .unwrap() + .contains("summary_cache_hit") + ); + let cached = response_json(response).await; + assert_eq!(cached["snapshot"]["total"], 1); + + let response = app + .clone() + .oneshot( + axum::http::Request::builder() + .method("POST") + .uri("/api/catalog") + .body(axum::body::Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + let (status, refreshed) = get_json(&app, uri).await; + assert_eq!(status, StatusCode::OK, "{refreshed}"); + assert_eq!(refreshed["snapshot"]["total"], 2); +} + #[tokio::test] async fn explorer_runs_prunes_unrelated_sources_before_querying() { let root = tempfile::tempdir().unwrap(); diff --git a/crates/persisting-pchronicle-cli/src/server/ui_cache.rs b/crates/persisting-pchronicle-cli/src/server/ui_cache.rs index f6b8a4829..010535baf 100644 --- a/crates/persisting-pchronicle-cli/src/server/ui_cache.rs +++ b/crates/persisting-pchronicle-cli/src/server/ui_cache.rs @@ -1,24 +1,18 @@ //! Rebuildable browse index. Accurate queries never use this index to establish //! source membership or revisions: an old index may omit newly created sources. -use std::collections::{HashMap, HashSet}; +use std::collections::{HashMap, HashSet, VecDeque}; use std::path::PathBuf; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; use anyhow::{Context, Result}; -use futures::TryStreamExt; -use lance::Dataset; -use lance::dataset::{InsertBuilder, MergeInsertBuilder, WhenMatched, WhenNotMatched}; -use lance::deps::arrow_array::{Array, RecordBatch, RecordBatchIterator, StringArray}; -use lance::deps::arrow_schema::{DataType, Field, Schema}; -use persisting_pchronicle::storage::{DatasetLocation, DatasetMount}; +use persisting_pchronicle::storage::{DatasetLocation, DatasetMount, ManifestCache}; use serde::{Deserialize, Serialize}; use tokio::sync::{RwLock, mpsc, oneshot}; use super::explorer::{CatalogTree, catalog_tree_from_mount_specs, catalog_tree_from_path_list}; -const CACHE_SCHEMA_VERSION: &str = "ui-tree-v2"; const REFRESH_INTERVAL: Duration = Duration::from_secs(30); const QUEUE_CAPACITY: usize = 128; // ponytail: one process-wide browse scan at a time; per-backend budgets if @@ -115,7 +109,7 @@ type Pending = Arc>>>; /// The task owns the index, not the coordinator; dropping the last AppState /// aborts it, including an in-flight list. No permanent process singleton. pub(crate) struct BrowseCoordinator { - index: Arc, + index: Arc, pending: Pending, refresh: Arc>>, sender: mpsc::Sender, @@ -145,14 +139,18 @@ impl BrowseCoordinator { .collect(); identities.sort(); let namespace = blake3::hash(&serde_json::to_vec(&identities).unwrap()).to_hex(); - let index = - Arc::new(CatalogIndex::open(root.join(format!("catalog-{namespace}.lance"))).await); + let index = Arc::new( + BrowseTreeProjection::open(root.join(format!("catalog-{namespace}.lance"))).await, + ); + let manifests = + Arc::new(ManifestCache::open(root.join(format!("manifest-{namespace}.lance"))).await); let pending = Arc::new(Mutex::new(HashMap::new())); let refresh = Arc::new(Mutex::new(HashMap::new())); let (sender, receiver) = mpsc::channel(QUEUE_CAPACITY); let task = tokio::spawn(run_worker( mounts, index.clone(), + manifests.clone(), pending.clone(), refresh.clone(), receiver, @@ -166,6 +164,31 @@ impl BrowseCoordinator { } } + pub(crate) async fn cached_dataset(&self, mount: &DatasetMount) -> Option { + let key = TreeKey::new(mount, "").ok()?; + self.index + .values + .read() + .await + .get(&key) + .map(|entry| entry.tree.clone()) + } + + pub(crate) async fn cached_source_paths(&self, mount: &DatasetMount) -> Vec { + let values = self.index.values.read().await; + let fingerprint = mount_fingerprint(mount); + let mut paths = values + .iter() + .filter(|(key, _)| key.dataset == mount.name && key.uri_fingerprint == fingerprint) + .flat_map(|(_, entry)| entry.tree.children.iter()) + .filter(|child| child.kind == "file") + .map(|child| child.path.clone()) + .collect::>(); + paths.sort(); + paths.dedup(); + paths + } + pub(crate) async fn roots(&self, mounts: &[DatasetMount]) -> BrowseSnapshot { let mut tree = catalog_tree_from_mount_specs(mounts); let values = self.index.values.read().await; @@ -296,7 +319,8 @@ impl BrowseCoordinator { async fn run_worker( mounts: Vec, - index: Arc, + index: Arc, + manifests: Arc, pending: Pending, states: Arc>>, mut receiver: mpsc::Receiver, @@ -305,7 +329,9 @@ async fn run_worker( interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); // Incremental round over roots and previously browsed prefixes. Foreground // requests are checked between each scan, not after the entire mount set. - let mut background = Vec::new(); + // Breadth-first keeps shallow datasets visible while a deep subtree is + // still being indexed. + let mut background = VecDeque::new(); let mut visited = HashSet::new(); loop { let key = tokio::select! { @@ -319,7 +345,7 @@ async fn run_worker( } continue; } - _ = std::future::ready(()), if !background.is_empty() => background.pop().unwrap(), + _ = std::future::ready(()), if !background.is_empty() => background.pop_front().unwrap(), }; let Some(mount) = mounts .iter() @@ -351,9 +377,14 @@ async fn run_worker( .await .context("browse path unavailable")?; } - let entries = tokio::time::timeout(Duration::from_secs(20), location.list(&key.prefix)) - .await - .context("browse list timed out")??; + let manifest_key = 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), + ) + .await + .context("browse list timed out")??; + let entries = listing.entries; let tree = catalog_tree_from_path_list(&mount.name, &key.prefix, &entries); // Walk only navigational directories; a Dataset leaf is opaque. // A bounded frontier prevents the background walk growing without limit. @@ -362,7 +393,7 @@ async fn run_worker( if child.kind == "dir" && visited.len() < 10_000 { let next = TreeKey::new(mount, &child.path)?; if visited.insert(next.clone()) { - background.push(next); + background.push_back(next); } } } @@ -383,6 +414,7 @@ async fn run_worker( } Err(error) => { let error = format!("{error:#}"); + let cached_view = index.values.read().await.contains_key(&key); let mut states = states.lock().unwrap(); let state = states.entry(key.clone()).or_default(); state.failures = state.failures.saturating_add(1); @@ -391,7 +423,7 @@ async fn run_worker( + Duration::from_secs((30u64 * (1u64 << state.failures.min(4))).min(300)), ); state.error = Some(error.clone()); - tracing::warn!(target: "pchronicle.serve", dataset = %key.dataset, prefix = %key.prefix, error = %error, "browse refresh failed; retaining previous view"); + tracing::warn!(target: "pchronicle.serve", dataset = %key.dataset, prefix = %key.prefix, cached_view, error = %error, "browse refresh failed"); Err(error) } }; @@ -411,52 +443,30 @@ fn now() -> i64 { chrono::Utc::now().timestamp() } -/// Disk is optional: corruption or a competing process never prevents browsing. -struct CatalogIndex { - path: PathBuf, +/// UI-only tree projection over the core ManifestCache. It preserves the legacy tree wire format. +struct BrowseTreeProjection { + disk: Arc, values: RwLock>, - // Hold an advisory lock for the writer lifetime, including corruption repair. - // A second process uses memory only instead of deleting/writing active data. - _disk_lock: Option, } -impl CatalogIndex { +impl BrowseTreeProjection { async fn open(path: PathBuf) -> Self { - let lock = (|| -> Result { - std::fs::create_dir_all(path.parent().context("cache parent")?)?; - let file = std::fs::OpenOptions::new() - .create(true) - .truncate(false) - .read(true) - .write(true) - .open(path.with_extension("lock"))?; - fs2::FileExt::try_lock_exclusive(&file)?; - Ok(file) - })(); - let mut index = Self { - path, - values: RwLock::new(HashMap::new()), - _disk_lock: lock.ok(), - }; - if index._disk_lock.is_none() { - return index; - } - match index.load().await { - Ok(values) => *index.values.get_mut() = values, - Err(error) => { - tracing::warn!(target: "pchronicle.serve", error = %error, "browse cache unreadable; rebuilding"); - let removed = if index.path.is_dir() { - tokio::fs::remove_dir_all(&index.path).await - } else { - tokio::fs::remove_file(&index.path).await - }; - if removed.is_err() && index.path.exists() { - // No repeated broken-dataset writes if repair is impossible. - index._disk_lock = None; - } - } + let disk = Arc::new(ManifestCache::open(path).await); + let values = disk + .projection_values() + .await + .into_iter() + .filter_map(|(key, value)| { + Some(( + serde_json::from_str(&key).ok()?, + serde_json::from_value(value).ok()?, + )) + }) + .collect(); + Self { + disk, + values: RwLock::new(values), } - index } async fn put(&self, key: TreeKey, tree: CatalogTree) { @@ -502,104 +512,35 @@ impl CatalogIndex { // Unchanged directory observations update memory without creating a // new Lance version every 30 seconds. On restart the older timestamp // conservatively marks the persisted view stale until revalidated. - if self._disk_lock.is_some() + if self.disk.writable() && (changed || !removed.is_empty()) - && let Err(error) = self.persist(&key, &entry, &removed).await - { - tracing::warn!(target: "pchronicle.serve", error = %error, "browse cache persistence failed; using memory"); - } - } - - async fn load(&self) -> Result> { - if !self.path.exists() { - return Ok(HashMap::new()); - } - let dataset = Dataset::open(self.path.to_string_lossy().as_ref()).await?; - let batches: Vec = dataset - .scan() - .try_into_stream() - .await? - .try_collect() - .await?; - let mut values = HashMap::new(); - for batch in batches { - let keys = text_column(&batch, "key")?; - let payloads = text_column(&batch, "payload")?; - let versions = text_column(&batch, "schema_version")?; - for row in 0..batch.num_rows() { - anyhow::ensure!( - versions[row] == CACHE_SCHEMA_VERSION, - "unsupported browse cache schema" - ); - values.insert( - serde_json::from_str(&keys[row])?, - serde_json::from_str(&payloads[row])?, - ); - } - } - Ok(values) - } - - async fn persist(&self, key: &TreeKey, entry: &IndexEntry, removed: &[TreeKey]) -> Result<()> { - let schema = Arc::new(Schema::new(vec![ - Field::new("key", DataType::Utf8, false), - Field::new("payload", DataType::Utf8, false), - Field::new("schema_version", DataType::Utf8, false), - ])); - let batch = RecordBatch::try_new( - schema.clone(), - vec![ - Arc::new(StringArray::from(vec![serde_json::to_string(key)?])) as _, - Arc::new(StringArray::from(vec![serde_json::to_string(entry)?])) as _, - Arc::new(StringArray::from(vec![CACHE_SCHEMA_VERSION])) as _, - ], - )?; - if self.path.exists() { - let mut dataset = Dataset::open(self.path.to_string_lossy().as_ref()).await?; - for chunk in removed.chunks(128) { - let keys = chunk + && let Err(error) = async { + let removed_keys = removed .iter() - .map(|key| { - Ok(format!( - "'{}'", - serde_json::to_string(key)?.replace("'", "''") - )) - }) - .collect::>>()? - .join(","); - dataset.delete(&format!("key IN ({keys})")).await?; + .map(|key| serde_json::to_string(key).unwrap()) + .collect::>(); + self.disk.remove_projections(&removed_keys).await?; + self.disk + .put_projection( + serde_json::to_string(&key).unwrap(), + &serde_json::to_value(&entry).unwrap(), + ) + .await } - let reader = Box::new(RecordBatchIterator::new(vec![Ok(batch)], schema)); - MergeInsertBuilder::try_new(Arc::new(dataset), vec!["key".into()])? - .when_matched(WhenMatched::UpdateAll) - .when_not_matched(WhenNotMatched::InsertAll) - .try_build()? - .execute_reader(reader) - .await?; - } else { - InsertBuilder::new(self.path.to_string_lossy().as_ref()) - .execute(vec![batch]) - .await?; + .await + { + tracing::warn!(target: "pchronicle.serve", error = %error, "browse cache persistence failed; using memory"); } - Ok(()) } } -fn text_column(batch: &RecordBatch, name: &str) -> Result> { - let array = batch - .column(batch.schema().index_of(name)?) - .as_any() - .downcast_ref::() - .with_context(|| format!("browse cache column {name} must be Utf8"))?; - anyhow::ensure!(array.null_count() == 0, "null browse cache field"); - Ok((0..array.len()) - .map(|index| array.value(index).to_owned()) - .collect()) -} - #[cfg(test)] mod tests { use super::*; + use lance::Dataset; + use lance::dataset::InsertBuilder; + use lance::deps::arrow_array::{RecordBatch, StringArray}; + use lance::deps::arrow_schema::{DataType, Field, Schema}; fn mount(path: &std::path::Path) -> DatasetMount { DatasetMount::new("test", path.to_string_lossy()).unwrap() @@ -610,13 +551,13 @@ mod tests { let temp = tempfile::tempdir().unwrap(); let path = temp.path().join("cache.lance"); std::fs::write(&path, "broken cache").unwrap(); - let index = CatalogIndex::open(path.clone()).await; + let index = BrowseTreeProjection::open(path.clone()).await; let key = TreeKey::new(&mount(temp.path()), "").unwrap(); index.put(key.clone(), CatalogTree::default()).await; assert!(path.is_dir()); assert_eq!(index.values.read().await.len(), 1); drop(index); - let reloaded = CatalogIndex::open(path).await; + let reloaded = BrowseTreeProjection::open(path).await; assert!(reloaded.values.read().await.contains_key(&key)); } @@ -635,7 +576,7 @@ mod tests { .execute(vec![batch]) .await .unwrap(); - let index = CatalogIndex::open(path.clone()).await; + let index = BrowseTreeProjection::open(path.clone()).await; assert!(index.values.read().await.is_empty()); assert!(!path.exists()); index @@ -656,7 +597,7 @@ mod tests { let second = BrowseCoordinator::start_at(vec![mount(&temp.path().join("b"))], temp.path().into()) .await; - assert_ne!(first.index.path, second.index.path); + assert_ne!(first.index.disk.path(), second.index.disk.path()); let weak = Arc::downgrade(&first.index); drop(first); tokio::time::timeout(Duration::from_secs(2), async { @@ -706,7 +647,7 @@ mod tests { async fn successful_parent_refresh_removes_deleted_descendants_on_disk() { let temp = tempfile::tempdir().unwrap(); let path = temp.path().join("cache.lance"); - let index = CatalogIndex::open(path.clone()).await; + let index = BrowseTreeProjection::open(path.clone()).await; let mount = mount(temp.path()); let child = TreeKey::new(&mount, "gone/nested").unwrap(); index.put(child.clone(), CatalogTree::default()).await; @@ -715,7 +656,7 @@ mod tests { .await; assert!(!index.values.read().await.contains_key(&child)); drop(index); - let reloaded = CatalogIndex::open(path).await; + let reloaded = BrowseTreeProjection::open(path).await; assert!(!reloaded.values.read().await.contains_key(&child)); } @@ -741,7 +682,7 @@ mod tests { async fn unchanged_observations_do_not_create_lance_versions() { let temp = tempfile::tempdir().unwrap(); let path = temp.path().join("cache.lance"); - let index = CatalogIndex::open(path.clone()).await; + let index = BrowseTreeProjection::open(path.clone()).await; let key = TreeKey::new(&mount(temp.path()), "").unwrap(); index.put(key.clone(), CatalogTree::default()).await; let version = Dataset::open(path.to_string_lossy().as_ref()) @@ -764,10 +705,10 @@ mod tests { async fn second_writer_does_not_remove_or_mutate_owned_cache() { let temp = tempfile::tempdir().unwrap(); let path = temp.path().join("cache.lance"); - let first = CatalogIndex::open(path.clone()).await; - let second = CatalogIndex::open(path.clone()).await; - assert!(first._disk_lock.is_some()); - assert!(second._disk_lock.is_none()); + let first = BrowseTreeProjection::open(path.clone()).await; + let second = BrowseTreeProjection::open(path.clone()).await; + assert!(first.disk.writable()); + assert!(!second.disk.writable()); second .put( TreeKey::new(&mount(temp.path()), "").unwrap(), diff --git a/crates/persisting-pchronicle/src/storage.rs b/crates/persisting-pchronicle/src/storage.rs index 41d25654e..1b23de4dd 100644 --- a/crates/persisting-pchronicle/src/storage.rs +++ b/crates/persisting-pchronicle/src/storage.rs @@ -64,29 +64,31 @@ pub use crate::store::object_store_io_gate::{ #[cfg(feature = "lance-store")] pub use crate::store::{ - AppendOutcome, AttemptRecord, AttemptRecordState, AttemptRegistry, CatalogDataset, - CatalogErrorPolicy, CatalogEventProvenance, CatalogEventView, CatalogNamespace, CatalogPage, - CatalogProjectionStatus, CatalogSnapshotOptions, CatalogSourceDescription, CatalogSourceKind, - CatalogSourceRevision, CatalogSourceStatus, CatalogStorylineKey, CatalogTrajectoryBundle, - ChronicleManifest, CommitRunOutcome, CompactJsonlBuildPhase, CompactJsonlColumn, - CompactJsonlImportEvent, CompactJsonlOffload, CompactJsonlOptions, CompactJsonlRecord, - CompactJsonlStore, DEFAULT_CONTENT_OFFLOAD_THRESHOLD, DEFAULT_CONTENT_PREVIEW_BYTES, - DEFAULT_DATASET_NAME, DEFAULT_MAX_CHUNK_BYTES, DEFAULT_MAX_EVENT_FALLBACK_BYTES, - DEFAULT_MAX_EVENT_FALLBACK_ROWS, DEFAULT_PHYSICAL_PAGE_LIMIT, DatasetCatalogSnapshot, - DatasetLocation, DatasetLocationKind, DatasetMount, DiscoveredSource, EventFactSnapshot, - EventLogLayoutStats, EventWriterFence, ExportOutcome, ImportableObjectEvent, - LanceMaintenanceOptions, LanceMaintenanceReport, LeaseAcquireOutcome, ManifestKind, - ManifestStats, NamespacePath, ObjectStoreManifestWriteMode, PathListEntry, PathListKind, - PhysicalColumn, PhysicalDataFile, PhysicalFileLayout, PhysicalFragment, PhysicalLayout, - PhysicalPage, PhysicalPagePreview, PhysicalPageQuery, PhysicalSource, PhysicalTable, - ProjectionSourceSnapshot, QueryScope, RawEventLanceAppender, RawEventLanceStore, ReplayOutcome, - RunControlStore, ShallowNavEntry, StorylineContentOptions, StorylineContentReadMode, - StorylineDataSource, StorylineDataSourceOptions, StorylineLanceStore, - StorylineMaintenanceReport, StorylineProjectionLineage, StorylineSearchIndexSuppressGuard, - StorylineStreamImportReport, StorylineStreamOptions, StorylineTablePaths, TrajectoryStats, - attempt_registry_now_ms, distinct_session_ids_in_run, export_source_dirs, export_story_bundle, - inspect_physical_file, inspect_physical_layout, inspect_physical_page, list_physical_sources, - load_manifest, load_manifest_at_uri, raw_event_lance_path, write_compact_jsonl_manifest, + AppendOutcome, AttemptRecord, AttemptRecordState, AttemptRegistry, CachedDataset, + CatalogDataset, CatalogErrorPolicy, CatalogEventProvenance, CatalogEventView, CatalogNamespace, + CatalogPage, CatalogProjectionStatus, CatalogSnapshotOptions, CatalogSourceDescription, + CatalogSourceKind, CatalogSourceRevision, CatalogSourceStatus, CatalogStorylineKey, + CatalogTrajectoryBundle, ChronicleManifest, CommitRunOutcome, CompactJsonlBuildPhase, + CompactJsonlColumn, CompactJsonlImportEvent, CompactJsonlOffload, CompactJsonlOptions, + CompactJsonlRecord, CompactJsonlStore, DEFAULT_CONTENT_OFFLOAD_THRESHOLD, + DEFAULT_CONTENT_PREVIEW_BYTES, DEFAULT_DATASET_NAME, DEFAULT_MAX_CHUNK_BYTES, + DEFAULT_MAX_EVENT_FALLBACK_BYTES, DEFAULT_MAX_EVENT_FALLBACK_ROWS, DEFAULT_PHYSICAL_PAGE_LIMIT, + Dataset, DatasetCatalogSnapshot, DatasetLocation, DatasetLocationKind, DatasetMount, + DatasetResolver, DiscoveredSource, EventFactSnapshot, EventLogLayoutStats, EventWriterFence, + ExportOutcome, ImportableObjectEvent, LanceMaintenanceOptions, LanceMaintenanceReport, + LeaseAcquireOutcome, LocationSummary, ManifestCache, ManifestKind, ManifestListing, + ManifestReadMode, ManifestStats, NamespacePath, ObjectStoreManifestWriteMode, PathListEntry, + PathListKind, PersistentCache, PhysicalColumn, PhysicalDataFile, PhysicalFileLayout, + PhysicalFragment, PhysicalLayout, PhysicalPage, PhysicalPagePreview, PhysicalPageQuery, + PhysicalSource, PhysicalTable, ProjectionSourceSnapshot, QueryScope, RawEventLanceAppender, + RawEventLanceStore, ReplayOutcome, ResolveMode, ResolveTarget, RunControlStore, + ShallowNavEntry, StorylineContentOptions, StorylineContentReadMode, StorylineDataSource, + StorylineDataSourceOptions, StorylineLanceStore, StorylineMaintenanceReport, + StorylineProjectionLineage, StorylineSearchIndexSuppressGuard, StorylineStreamImportReport, + StorylineStreamOptions, StorylineTablePaths, TrajectoryStats, attempt_registry_now_ms, + distinct_session_ids_in_run, export_source_dirs, export_story_bundle, inspect_physical_file, + inspect_physical_layout, inspect_physical_page, list_physical_sources, load_manifest, + load_manifest_at_uri, raw_event_lance_path, write_compact_jsonl_manifest, write_storyline_manifest, write_storyline_manifest_at_uri, }; diff --git a/crates/persisting-pchronicle/src/store/catalog/discovery.rs b/crates/persisting-pchronicle/src/store/catalog/discovery.rs index 239cd3b78..8ce877df1 100644 --- a/crates/persisting-pchronicle/src/store/catalog/discovery.rs +++ b/crates/persisting-pchronicle/src/store/catalog/discovery.rs @@ -1,5 +1,5 @@ use super::*; -use crate::store::chronicle_manifest::{ManifestKind, try_load_manifest}; +use crate::store::catalog::manifest::{ManifestKind, try_load_manifest}; use crate::store::opendal_store::Store as OpendalStore; #[derive(Debug)] @@ -140,7 +140,7 @@ pub(super) async fn freeze_candidate( generation: paths.generation.clone(), }); if let Ok(Some(manifest)) = - crate::store::chronicle_manifest::load_manifest_at_uri(&uri).await + crate::store::catalog::manifest::load_manifest_at_uri(&uri).await && manifest.is_storyline_leaf() && let Some(stats) = manifest.stats { @@ -184,7 +184,7 @@ pub(super) async fn freeze_candidate( } Candidate::Compact { file, uri, .. } => { if let Some(manifest) = - crate::store::chronicle_manifest::try_load_manifest(Path::new(&uri)) + crate::store::catalog::manifest::try_load_manifest(Path::new(&uri)) && let Some(stats) = manifest.stats { source_row.record_count = Some(stats.record_count); @@ -520,6 +520,74 @@ pub(super) async fn discover_candidates( } } +pub(super) async fn discover_cached_candidates( + mount: &DatasetMount, + files: &[String], + options: LocalQueryManifestOptions, +) -> Result> { + if local_mount_path(&mount.uri).is_none() { + // Remote discovery is mount-wide; do it once and filter the result. + let wanted: BTreeSet<_> = files.iter().map(|file| file.trim_matches('/')).collect(); + return Ok(discover_candidates(mount, options) + .await? + .into_iter() + .filter(|candidate| wanted.contains(candidate.source_stub().file.as_str())) + .collect()); + } + let mut candidates = Vec::new(); + for file in files { + candidates.extend(discover_cached_candidate_at(mount, file, options).await?); + } + Ok(candidates) +} + +pub(super) async fn discover_cached_candidate_at( + mount: &DatasetMount, + file: &str, + options: LocalQueryManifestOptions, +) -> Result> { + let file = file.trim().trim_matches('/'); + anyhow::ensure!( + !file.is_empty() + && !file + .split('/') + .any(|part| part.is_empty() || matches!(part, "." | "..") || part.contains('\\')), + "invalid cached source path" + ); + let Some(root) = local_mount_path(&mount.uri) else { + // Remote cache entries are hints only until object metadata can be + // validated without a full listing. + return discover_candidates(mount, options).await; + }; + let path = root.join(file); + let metadata = match fs::symlink_metadata(&path) { + Ok(metadata) => metadata, + Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()), + Err(error) => return Err(error).context("inspect cached source"), + }; + if metadata.is_file() { + anyhow::ensure!( + is_json_candidate(&path), + "cached source is not a supported JSON file" + ); + return Ok(vec![Candidate::LocalFile { + file: file.to_owned(), + root: root.clone(), + path, + size_bytes: metadata.len(), + last_modified: modified_string(&metadata), + }]); + } + let candidates = match classify_local_dir(&root, &path).await? { + LocalDirClass::Leaf(candidates) => candidates, + LocalDirClass::Skip => Vec::new(), + LocalDirClass::Recurse => { + collect_local_virtual(&root, &path, DiscoveryBudget::new(options)).await? + } + }; + Ok(candidates) +} + pub(super) async fn discover_candidate_at( mount: &DatasetMount, file: &str, diff --git a/crates/persisting-pchronicle/src/store/location.rs b/crates/persisting-pchronicle/src/store/catalog/location.rs similarity index 91% rename from crates/persisting-pchronicle/src/store/location.rs rename to crates/persisting-pchronicle/src/store/catalog/location.rs index ae544c1bf..e0e929ff7 100644 --- a/crates/persisting-pchronicle/src/store/location.rs +++ b/crates/persisting-pchronicle/src/store/catalog/location.rs @@ -6,9 +6,10 @@ use std::io::{ErrorKind, Write}; use std::path::{Path, PathBuf}; use anyhow::{Context, Result, anyhow}; +use futures::{StreamExt, TryStreamExt}; use url::Url; -use super::opendal_store::Store as OpendalStore; +use crate::store::opendal_store::Store as OpendalStore; /// One discovery event while walking importable JSON objects. #[derive(Debug, Clone)] @@ -34,7 +35,7 @@ pub struct ShallowNavEntry { } /// One child of RFC-0015 `list(path)` — shell-like one-level listing. -#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] #[serde(rename_all = "snake_case")] pub enum PathListKind { Directory, @@ -42,7 +43,7 @@ pub enum PathListKind { File, } -#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)] +#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] pub struct PathListEntry { pub name: String, pub path: String, @@ -281,7 +282,7 @@ impl DatasetLocation { if !dir.is_dir() { return Ok(None); } - if let Some(manifest) = crate::store::chronicle_manifest::try_load_manifest(&dir) { + if let Some(manifest) = crate::store::catalog::manifest::try_load_manifest(&dir) { if manifest.is_storyline_leaf() { return Ok(Some("storyline")); } @@ -357,7 +358,8 @@ impl DatasetLocation { !relative.split('/').any(|part| part == ".."), "relative object path must not contain '..'" ); - if let Some(preview) = self.list_dataset_preview(relative).await? { + if let Some(kind) = self.probe_nav_dataset_kind(relative).await? { + let preview = self.list_dataset_preview(relative, kind).await?; let name = if relative.is_empty() { ".".to_string() } else { @@ -377,7 +379,7 @@ impl DatasetLocation { }]); } - let nav = self.list_shallow_nav(relative).await?; + let nav = self.list_nav_children(relative).await?; let mut out = Vec::with_capacity(nav.len()); for entry in nav { let path = if relative.is_empty() { @@ -385,7 +387,8 @@ impl DatasetLocation { } else { format!("{relative}/{}", entry.name) }; - if let Some(preview) = self.list_dataset_preview(&path).await? { + if let Some(kind) = entry.dataset_kind.as_deref() { + let preview = self.list_dataset_preview(&path, kind).await?; out.push(PathListEntry { name: entry.name, path, @@ -417,10 +420,7 @@ impl DatasetLocation { Ok(out) } - async fn list_dataset_preview(&self, relative: &str) -> Result> { - let Some(kind) = self.probe_nav_dataset_kind(relative).await? else { - return Ok(None); - }; + async fn list_dataset_preview(&self, relative: &str, kind: &str) -> Result { let mut preview = DatasetListPreview { format: Some(match kind { "storyline" => "storyline-lance".into(), @@ -436,7 +436,7 @@ impl DatasetLocation { } else { root.join(relative) }; - if let Some(manifest) = crate::store::chronicle_manifest::try_load_manifest(&dir) + if let Some(manifest) = crate::store::catalog::manifest::try_load_manifest(&dir) && matches!(manifest.kind, crate::store::ManifestKind::Leaf) { if let Some(format) = manifest.format { @@ -447,7 +447,7 @@ impl DatasetLocation { preview.failed_count = Some(stats.failed_count); } } - return Ok(Some(preview)); + return Ok(preview); } let store = OpendalStore::from_uri(&self.uri).await?; let key = if relative.is_empty() { @@ -470,7 +470,7 @@ impl DatasetLocation { preview.failed_count = Some(stats.failed_count); } } - Ok(Some(preview)) + Ok(preview) } /// Immediate children under a Dataset-relative prefix for explorer navigation. @@ -491,6 +491,11 @@ impl DatasetLocation { if self.probe_nav_dataset_kind(relative).await?.is_some() { return Ok(Vec::new()); } + self.list_nav_children(relative).await + } + + // Caller has validated the path and established that it is not a Dataset leaf. + async fn list_nav_children(&self, relative: &str) -> Result> { if let Some(root) = &self.local_path { let dir = if relative.is_empty() { root.clone() @@ -586,27 +591,24 @@ impl DatasetLocation { } dirs.insert(child.to_string()); } - let mut out = Vec::with_capacity(dirs.len() + files.len()); - for name in dirs { + // Bound remote marker probes so wide directories do not serialize every + // round trip, without launching an unbounded number of storage requests. + let mut out: Vec<_> = futures::stream::iter(dirs.into_iter().map(|name| async move { let child_rel = if relative.is_empty() { name.clone() } else { format!("{relative}/{name}") }; - if let Some(kind) = self.probe_nav_dataset_kind(&child_rel).await? { - out.push(ShallowNavEntry { - name, - is_dir: false, - dataset_kind: Some(kind.into()), - }); - } else { - out.push(ShallowNavEntry { - name, - is_dir: true, - dataset_kind: None, - }); - } - } + let kind = self.probe_nav_dataset_kind(&child_rel).await?; + Ok::<_, anyhow::Error>(ShallowNavEntry { + name, + is_dir: kind.is_none(), + dataset_kind: kind.map(str::to_owned), + }) + })) + .buffered(8) + .try_collect() + .await?; for name in files { out.push(ShallowNavEntry { name, @@ -1228,46 +1230,78 @@ mod tests { let archive = warehouse.join("archive"); std::fs::create_dir_all(team.join("codex_jsonl")).unwrap(); std::fs::create_dir_all(&archive).unwrap(); - crate::store::chronicle_manifest::write_compact_jsonl_manifest(&archive, 1, 7).unwrap(); + crate::store::catalog::manifest::write_compact_jsonl_manifest(&archive, 1, 7).unwrap(); std::fs::write(warehouse.join("notes.json"), b"[]").unwrap(); - let location = DatasetLocation::parse(warehouse.to_str().unwrap()).unwrap(); - let listed = location.list("").await.unwrap(); - let summary: Vec<_> = listed - .iter() - .map(|entry| { - ( - entry.name.as_str(), - entry.kind, - entry.record_count, - entry.format.as_deref(), - ) - }) - .collect(); - assert_eq!( - summary, - vec![ - ( - "archive", - PathListKind::Dataset, - Some(7), - Some("compact-jsonl/v1") - ), - ("notes.json", PathListKind::File, None, None), - ("team", PathListKind::Directory, None, None), - ] - ); - assert!( - !listed + let remote = DatasetLocation::parse(&format!( + "shared-memory://pchronicle-list-{}/warehouse", + uuid::Uuid::new_v4().simple() + )) + .unwrap(); + remote + .write_relative_bytes("team/codex_jsonl/data.jsonl", b"{}") + .await + .unwrap(); + remote + .write_relative_bytes("notes.json", b"[]") + .await + .unwrap(); + remote + .write_relative_bytes( + "archive/chronicle.manifest", + &std::fs::read(archive.join("chronicle.manifest")).unwrap(), + ) + .await + .unwrap(); + + for location in [ + DatasetLocation::parse(warehouse.to_str().unwrap()).unwrap(), + remote, + ] { + let listed = location.list("").await.unwrap(); + let summary: Vec<_> = listed .iter() - .any(|entry| entry.path.contains("codex_jsonl")) - ); + .map(|entry| { + ( + entry.name.as_str(), + entry.kind, + entry.record_count, + entry.format.as_deref(), + ) + }) + .collect(); + assert_eq!( + summary, + vec![ + ( + "archive", + PathListKind::Dataset, + Some(7), + Some("compact-jsonl/v1") + ), + ("notes.json", PathListKind::File, None, None), + ("team", PathListKind::Directory, None, None), + ] + ); + assert!( + !listed + .iter() + .any(|entry| entry.path.contains("codex_jsonl")) + ); - let leaf = location.list("archive").await.unwrap(); - assert_eq!(leaf.len(), 1); - assert_eq!(leaf[0].kind, PathListKind::Dataset); - assert_eq!(leaf[0].record_count, Some(7)); - assert_eq!(leaf[0].path, "archive"); + let leaf = location.list("archive").await.unwrap(); + assert_eq!(leaf.len(), 1); + assert_eq!(leaf[0].kind, PathListKind::Dataset); + assert_eq!(leaf[0].record_count, Some(7)); + assert_eq!(leaf[0].path, "archive"); + assert!( + location + .list_shallow_nav("archive") + .await + .unwrap() + .is_empty() + ); + } } #[tokio::test] diff --git a/crates/persisting-pchronicle/src/store/chronicle_manifest.rs b/crates/persisting-pchronicle/src/store/catalog/manifest.rs similarity index 98% rename from crates/persisting-pchronicle/src/store/chronicle_manifest.rs rename to crates/persisting-pchronicle/src/store/catalog/manifest.rs index 6aafd1ce8..1987e9f99 100644 --- a/crates/persisting-pchronicle/src/store/chronicle_manifest.rs +++ b/crates/persisting-pchronicle/src/store/catalog/manifest.rs @@ -262,7 +262,7 @@ pub async fn write_storyline_manifest_at_uri( failed_count, ); manifest.validate()?; - let location = crate::store::location::DatasetLocation::parse(root_uri) + let location = crate::store::catalog::location::DatasetLocation::parse(root_uri) .with_context(|| format!("parse Dataset URI for chronicle.manifest ({root_uri})"))?; if let Some(path) = location.local_path() { return atomic_write_manifest(path, &manifest); @@ -276,7 +276,7 @@ pub async fn write_storyline_manifest_at_uri( /// Load a manifesto from a local path or object-store Dataset URI. pub async fn load_manifest_at_uri(root_uri: &str) -> Result> { - let location = crate::store::location::DatasetLocation::parse(root_uri) + let location = crate::store::catalog::location::DatasetLocation::parse(root_uri) .with_context(|| format!("parse Dataset URI for chronicle.manifest ({root_uri})"))?; if let Some(path) = location.local_path() { return load_manifest(path); diff --git a/crates/persisting-pchronicle/src/store/catalog/manifest_cache.rs b/crates/persisting-pchronicle/src/store/catalog/manifest_cache.rs new file mode 100644 index 000000000..6ecf404b5 --- /dev/null +++ b/crates/persisting-pchronicle/src/store/catalog/manifest_cache.rs @@ -0,0 +1,286 @@ +//! Rebuildable cache of manifest-derived directory observations. +//! +//! 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::path::PathBuf; +use std::sync::Arc; +use std::time::Duration; + +use anyhow::Result; +use serde::{Deserialize, Serialize}; +use tokio::sync::Mutex; +use tokio::sync::RwLock; + +use crate::store::{DatasetLocation, PathListEntry, PersistentCache}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ManifestListing { + pub entries: Vec, + pub observed_at: i64, +} + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct LocationSummary { + pub datasets: u64, + pub trajectories: u64, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ManifestReadMode { + Cached, + RefreshIfMissing, + Fresh, +} + +/// Persistent manifest cache. Keys are caller-owned stable identities, e.g. +/// `mount-uri\0relative-prefix`; this keeps the cache independent of UI types. +#[derive(Clone)] +pub struct ManifestCache { + disk: Arc>, + values: Arc>>, + refresh_gate: Arc>, +} + +impl ManifestCache { + pub async fn open(path: PathBuf) -> Self { + let disk = Arc::new(PersistentCache::open(path).await); + let values = Arc::new(RwLock::new( + disk.values() + .await + .into_iter() + .filter_map(|(key, value)| serde_json::from_value(value).ok().map(|v| (key, v))) + .collect(), + )); + Self { + disk, + values, + refresh_gate: Arc::new(Mutex::new(())), + } + } + + pub async fn get(&self, key: &str) -> Option { + self.values.read().await.get(key).cloned() + } + + pub async fn get_or_refresh( + &self, + key: impl Into + Clone, + location: &DatasetLocation, + prefix: &str, + mode: ManifestReadMode, + ) -> Result> { + let key = key.into(); + if !matches!(mode, ManifestReadMode::Fresh) { + if let Some(value) = self.get(&key).await { + return Ok(Some(value)); + } + if matches!(mode, ManifestReadMode::Cached) { + return Ok(None); + } + } + self.refresh(key, location, prefix).await.map(Some) + } + + /// Read one level from the authoritative location and publish it + /// atomically. Callers may use this for a foreground miss or a worker. + pub async fn refresh( + &self, + key: impl Into, + location: &DatasetLocation, + prefix: &str, + ) -> Result { + let _guard = self.refresh_gate.lock().await; + let key = key.into(); + let listing = ManifestListing { + entries: location.list(prefix).await?, + observed_at: chrono::Utc::now().timestamp(), + }; + self.values + .write() + .await + .insert(key.clone(), listing.clone()); + self.disk + .upsert(&key, &serde_json::to_value(&listing)?, &[]) + .await?; + Ok(listing) + } + + pub fn writable(&self) -> bool { + self.disk.writable() + } + + pub fn path(&self) -> &std::path::Path { + self.disk.path() + } + + pub async fn projection_values(&self) -> std::collections::HashMap { + self.disk + .values() + .await + .into_iter() + .filter(|(key, _)| serde_json::from_str::(key).is_ok()) + .collect() + } + + pub async fn put_projection( + &self, + key: impl Into, + value: &serde_json::Value, + ) -> Result<()> { + self.disk.upsert(&key.into(), value, &[]).await + } + + pub async fn remove_projections(&self, keys: &[String]) -> Result<()> { + self.disk + .upsert( + &"__projection_tombstone__".to_owned(), + &serde_json::Value::Null, + keys, + ) + .await + } + + /// 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 { + let values = self.values.read().await; + values + .iter() + .filter(|(key, _)| key == &key_prefix || key.starts_with(&format!("{key_prefix}\0"))) + .fold(LocationSummary::default(), |mut total, (_, listing)| { + for entry in &listing.entries { + if matches!(entry.kind, crate::store::PathListKind::Dataset) { + total.datasets += 1; + total.trajectories += entry.record_count.unwrap_or_default(); + } + } + total + }) + } + + /// Breadth-first refresh of a mount. Only one refresh runs at a time; + /// 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; + let mut queue = VecDeque::from([String::new()]); + while let Some(prefix) = queue.pop_front() { + let listing = ManifestListing { + entries: location.list(&prefix).await?, + observed_at: chrono::Utc::now().timestamp(), + }; + let key = if prefix.is_empty() { + key_prefix.to_owned() + } else { + format!("{key_prefix}\0{prefix}") + }; + self.values + .write() + .await + .insert(key.clone(), listing.clone()); + self.disk + .upsert(&key, &serde_json::to_value(&listing)?, &[]) + .await?; + for child in listing + .entries + .iter() + .filter(|e| matches!(e.kind, crate::store::PathListKind::Directory)) + { + queue.push_back(child.path.clone()); + } + tokio::task::yield_now().await; + } + Ok(()) + } + + pub fn spawn_periodic_refresh( + &self, + key_prefix: String, + location: DatasetLocation, + interval: Duration, + ) -> tokio::task::JoinHandle<()> { + let cache = self.clone(); + tokio::spawn(async move { + let mut ticker = tokio::time::interval(interval); + loop { + ticker.tick().await; + if let Err(error) = cache.refresh_mount(&key_prefix, &location).await { + tracing::warn!(target: "pchronicle.serve", %error, "manifest cache refresh failed"); + } + } + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::store::{PathListEntry, PathListKind}; + + #[tokio::test] + async fn summary_counts_cached_manifest_entries() { + let dir = tempfile::tempdir().unwrap(); + let cache = ManifestCache::open(dir.path().join("manifest.lance")).await; + cache.values.write().await.insert( + "mount".into(), + ManifestListing { + entries: vec![PathListEntry { + name: "a".into(), + path: "a".into(), + kind: PathListKind::Dataset, + format: None, + record_count: Some(3), + failed_count: None, + }], + observed_at: 0, + }, + ); + assert_eq!( + cache.summary("mount").await, + LocationSummary { + datasets: 1, + trajectories: 3 + } + ); + } + + #[tokio::test] + async fn read_mode_uses_cache_then_refreshes_and_survives_restart() { + let dir = tempfile::tempdir().unwrap(); + let dataset = dir.path().join("leaf"); + std::fs::create_dir_all(&dataset).unwrap(); + crate::store::catalog::manifest::write_compact_jsonl_manifest(&dataset, 1, 7).unwrap(); + let location = DatasetLocation::parse(dir.path().to_str().unwrap()).unwrap(); + let cache_path = dir.path().join("cache.lance"); + let cache = ManifestCache::open(cache_path.clone()).await; + assert!( + cache + .get_or_refresh("mount", &location, "", ManifestReadMode::Cached) + .await + .unwrap() + .is_none() + ); + let listing = cache + .get_or_refresh("mount", &location, "", ManifestReadMode::RefreshIfMissing) + .await + .unwrap() + .unwrap(); + assert_eq!(listing.entries.len(), 1); + assert!( + cache + .get_or_refresh("mount", &location, "", ManifestReadMode::Cached) + .await + .unwrap() + .is_some() + ); + drop(cache); + assert!( + ManifestCache::open(cache_path) + .await + .get("mount") + .await + .is_some() + ); + } +} diff --git a/crates/persisting-pchronicle/src/store/catalog/mod.rs b/crates/persisting-pchronicle/src/store/catalog/mod.rs index 7dccb75c3..61f101834 100644 --- a/crates/persisting-pchronicle/src/store/catalog/mod.rs +++ b/crates/persisting-pchronicle/src/store/catalog/mod.rs @@ -7,18 +7,37 @@ mod discovery; mod identity; +pub mod location; +pub mod manifest; +mod manifest_cache; mod namespace; mod provider; +mod resolver; mod source; pub use identity::{CatalogSourceRevision, DatasetMount, NamespacePath}; +#[allow(unused_imports)] +pub use location::{ + DatasetLocation, DatasetLocationKind, ImportableObjectEvent, PathListEntry, PathListKind, + ShallowNavEntry, +}; +#[allow(unused_imports)] +pub use manifest::{CHRONICLE_MANIFEST_FILE, ChronicleManifest, ManifestKind, ManifestStats}; +#[allow(unused_imports)] +pub use manifest::{ + STORYLINE_FORMAT, atomic_write_manifest, compact_jsonl_manifest_matches, load_manifest, + load_manifest_at_uri, try_load_manifest, write_compact_jsonl_manifest, + write_storyline_manifest, write_storyline_manifest_at_uri, +}; +pub use manifest_cache::{LocationSummary, ManifestCache, ManifestListing, ManifestReadMode}; pub use namespace::{CatalogNamespace, CatalogPage, CatalogSourceDescription}; use provider::*; +pub use resolver::{CachedDataset, Dataset, DatasetResolver, ResolveMode, ResolveTarget}; use source::*; use discovery::{ - bind_canonical_storyline_projections, discover_candidate_at, discover_candidates, - freeze_candidate, normalize_event_storylines, + bind_canonical_storyline_projections, discover_cached_candidates, discover_candidate_at, + discover_candidates, freeze_candidate, normalize_event_storylines, }; use std::collections::{BTreeMap, BTreeSet, HashSet}; @@ -340,7 +359,7 @@ impl DatasetCatalogSnapshot { default_dataset: Option, options: CatalogSnapshotOptions, ) -> Result { - Self::discover_impl(mounts, default_dataset, options, None).await + Self::discover_impl(mounts, default_dataset, options, None, None).await } pub async fn discover_scoped( @@ -349,7 +368,27 @@ impl DatasetCatalogSnapshot { options: CatalogSnapshotOptions, scope: QueryScope, ) -> Result { - Self::discover_impl(mounts, default_dataset, options, Some(scope)).await + Self::discover_impl(mounts, default_dataset, options, Some(scope), None).await + } + + /// Build a snapshot from source paths already observed by a UI cache. + /// Missing or unsupported cached paths are ignored; callers that require + /// exact membership must use [`Self::discover`] instead. + pub async fn discover_scoped_from_cached_files( + mounts: Vec, + default_dataset: Option, + options: CatalogSnapshotOptions, + scope: QueryScope, + cached_files: Vec, + ) -> Result { + Self::discover_impl( + mounts, + default_dataset, + options, + Some(scope), + Some(&cached_files), + ) + .await } async fn discover_impl( @@ -357,6 +396,7 @@ impl DatasetCatalogSnapshot { default_dataset: Option, options: CatalogSnapshotOptions, scope: Option, + cached_files: Option<&[String]>, ) -> Result { anyhow::ensure!(!mounts.is_empty(), "mount at least one Dataset"); validate_catalog_options(options)?; @@ -393,7 +433,12 @@ impl DatasetCatalogSnapshot { Some(scope) if scope.dataset != mount.name => Vec::new(), Some(scope) => match scope.source_file.as_deref() { Some(file) => discover_candidate_at(&mount, file, options.manifest).await?, - None => discover_candidates(&mount, options.manifest).await?, + None => match cached_files { + Some(files) => { + discover_cached_candidates(&mount, files, options.manifest).await? + } + None => discover_candidates(&mount, options.manifest).await?, + }, }, None => discover_candidates(&mount, options.manifest).await?, }; diff --git a/crates/persisting-pchronicle/src/store/catalog/resolver.rs b/crates/persisting-pchronicle/src/store/catalog/resolver.rs new file mode 100644 index 000000000..3a0f9a0be --- /dev/null +++ b/crates/persisting-pchronicle/src/store/catalog/resolver.rs @@ -0,0 +1,221 @@ +//! One resolver for atomic Datasets and recursive DatasetMounts. +//! +//! A Dataset is a single queryable leaf. A DatasetMount is a named path whose +//! recursive contents may contain many leaves. The mode is explicit so a UI +//! cache can never silently become an authoritative CLI resolution. +use super::{ + CatalogSnapshotOptions, CatalogSourceKind, DatasetCatalogSnapshot, DatasetMount, QueryScope, +}; +use crate::store::DatasetLocation; +use anyhow::{Result, bail}; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum ResolveMode { + /// Resolve only source paths supplied by the UI index. + Cached, + /// Discover and pin the requested scope from the backing location. + Fresh, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum ResolveTarget { + /// Resolve all configured mounts together in one query snapshot. + Catalog, + Dataset { + mount: String, + file: String, + }, + Mount { + mount: String, + prefix: Option, + }, +} + +/// One atomic, queryable dataset leaf. It is deliberately separate from a +/// DatasetMount, which is a recursive namespace and may contain many leaves. +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct Dataset { + location: DatasetLocation, +} + +/// A persisted, manifest-derived identity for one atomic dataset under a +/// mount. It contains no open handles and is safe to use as a UI hint. +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct CachedDataset { + mount: String, + file: String, +} + +impl CachedDataset { + pub fn new(mount: impl Into, file: impl Into) -> Self { + Self { + mount: mount.into(), + file: file.into(), + } + } + + pub fn mount(&self) -> &str { + &self.mount + } + pub fn file(&self) -> &str { + &self.file + } +} + +impl Dataset { + pub fn new(uri: impl AsRef) -> Result { + Ok(Self { + location: DatasetLocation::parse(uri.as_ref())?, + }) + } + + pub fn uri(&self) -> &str { + self.location.as_str() + } + + pub async fn resolve(&self, options: CatalogSnapshotOptions) -> Result { + let mount = DatasetMount::new("dataset", self.uri())?; + let snapshot = DatasetCatalogSnapshot::discover_scoped( + vec![mount], + Some("dataset".into()), + options, + QueryScope { + dataset: "dataset".into(), + source_file: Some(".".into()), + }, + ) + .await?; + anyhow::ensure!( + snapshot.datasets().iter().all(|dataset| { + dataset + .sources + .iter() + .all(|source| source.kind != CatalogSourceKind::Directory) + }), + "atomic Dataset resolved to a recursive DatasetMount" + ); + let leaves = snapshot + .datasets() + .iter() + .flat_map(|dataset| dataset.sources.iter()) + .filter(|source| source.kind != CatalogSourceKind::Directory) + .count(); + anyhow::ensure!( + leaves == 1, + "atomic Dataset must resolve to exactly one leaf" + ); + Ok(snapshot) + } +} + +#[derive(Clone, Debug)] +pub struct DatasetResolver { + mounts: Vec, + default_dataset: Option, + options: CatalogSnapshotOptions, +} + +impl DatasetResolver { + pub fn new( + mounts: Vec, + default_dataset: Option, + options: CatalogSnapshotOptions, + ) -> Self { + Self { + mounts, + default_dataset, + options, + } + } + + pub async fn resolve( + &self, + target: ResolveTarget, + mode: ResolveMode, + cached_datasets: &[CachedDataset], + ) -> Result { + let scope = match target { + ResolveTarget::Catalog => { + return match mode { + ResolveMode::Fresh => { + DatasetCatalogSnapshot::discover( + self.mounts.clone(), + self.default_dataset.clone(), + self.options, + ) + .await + } + ResolveMode::Cached => { + bail!("cached Catalog resolution requires a scoped mount") + } + }; + } + ResolveTarget::Dataset { mount, file } => QueryScope { + dataset: mount, + source_file: Some(file), + }, + ResolveTarget::Mount { mount, prefix } => QueryScope { + dataset: mount, + source_file: prefix, + }, + }; + let cached_files: Vec<_> = cached_datasets + .iter() + .filter(|dataset| dataset.mount == scope.dataset) + .map(|dataset| dataset.file.clone()) + .collect(); + match mode { + ResolveMode::Fresh => { + DatasetCatalogSnapshot::discover_scoped( + self.mounts.clone(), + self.default_dataset.clone(), + self.options, + scope, + ) + .await + } + ResolveMode::Cached if cached_files.is_empty() => { + bail!("cached Dataset resolution requires at least one indexed source") + } + ResolveMode::Cached => { + DatasetCatalogSnapshot::discover_scoped_from_cached_files( + self.mounts.clone(), + self.default_dataset.clone(), + self.options, + scope, + cached_files, + ) + .await + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn atomic_dataset_keeps_location_identity() { + let dataset = Dataset::new("memory://resolver-test").unwrap(); + assert_eq!(dataset.uri(), "memory://resolver-test"); + } + + #[tokio::test] + async fn cached_resolution_never_silently_discovers() { + let mount = DatasetMount::new("prod", "memory://resolver-test").unwrap(); + let resolver = DatasetResolver::new(vec![mount], None, CatalogSnapshotOptions::default()); + let error = resolver + .resolve( + ResolveTarget::Mount { + mount: "prod".into(), + prefix: None, + }, + ResolveMode::Cached, + &[], + ) + .await + .unwrap_err(); + assert!(error.to_string().contains("cached Dataset resolution")); + } +} diff --git a/crates/persisting-pchronicle/src/store/catalog/tests.rs b/crates/persisting-pchronicle/src/store/catalog/tests.rs index bc70872ba..340d0089c 100644 --- a/crates/persisting-pchronicle/src/store/catalog/tests.rs +++ b/crates/persisting-pchronicle/src/store/catalog/tests.rs @@ -293,7 +293,7 @@ async fn discovers_extensionless_compact_lance_dataset() -> Result<()> { assert_eq!(sources[0].file, "compact"); assert_eq!(sources[0].format.as_deref(), Some("compact-jsonl/v1")); - let manifest = crate::storage::load_manifest(&compact)?.expect("import writes manifesto"); + let manifest = super::manifest::load_manifest(&compact)?.expect("import writes manifest"); assert!(manifest.is_compact_jsonl_leaf()); assert_eq!(manifest.stats.as_ref().unwrap().record_count, 1); Ok(()) @@ -306,11 +306,11 @@ async fn discovers_nested_branch_and_leaf_chronicle_manifests_without_opening_la let warehouse = temp.path().join("warehouse"); let leaf = warehouse.join("codex_jsonl"); fs::create_dir_all(&leaf)?; - crate::store::chronicle_manifest::atomic_write_manifest( + crate::store::catalog::manifest::atomic_write_manifest( &warehouse, &crate::store::ChronicleManifest::branch(), )?; - crate::store::chronicle_manifest::write_compact_jsonl_manifest(&leaf, 1, 42)?; + crate::store::catalog::manifest::write_compact_jsonl_manifest(&leaf, 1, 42)?; // No Lance data/ tree: discovery must trust the leaf manifesto. let snapshot = DatasetCatalogSnapshot::discover( @@ -339,17 +339,17 @@ async fn discovers_multi_level_branch_tree_and_preserves_leaf_counts() -> Result for dir in [&warehouse, &team, &leaf_a, &leaf_b, &sibling] { fs::create_dir_all(dir)?; } - crate::store::chronicle_manifest::atomic_write_manifest( + crate::store::catalog::manifest::atomic_write_manifest( &warehouse, &crate::store::ChronicleManifest::branch(), )?; - crate::store::chronicle_manifest::atomic_write_manifest( + crate::store::catalog::manifest::atomic_write_manifest( &team, &crate::store::ChronicleManifest::branch(), )?; - crate::store::chronicle_manifest::write_compact_jsonl_manifest(&leaf_a, 1, 10)?; - crate::store::chronicle_manifest::write_compact_jsonl_manifest(&leaf_b, 2, 20)?; - crate::store::chronicle_manifest::write_compact_jsonl_manifest(&sibling, 3, 7)?; + crate::store::catalog::manifest::write_compact_jsonl_manifest(&leaf_a, 1, 10)?; + crate::store::catalog::manifest::write_compact_jsonl_manifest(&leaf_b, 2, 20)?; + crate::store::catalog::manifest::write_compact_jsonl_manifest(&sibling, 3, 7)?; let snapshot = DatasetCatalogSnapshot::discover( vec![DatasetMount::default(warehouse.to_string_lossy())?], @@ -385,8 +385,8 @@ async fn leaf_manifest_does_not_recurse_into_nested_child_manifest() -> Result<( let leaf = temp.path().join("leaf"); let nested = leaf.join("nested_child"); fs::create_dir_all(&nested)?; - crate::store::chronicle_manifest::write_compact_jsonl_manifest(&leaf, 1, 5)?; - crate::store::chronicle_manifest::write_compact_jsonl_manifest(&nested, 1, 99)?; + crate::store::catalog::manifest::write_compact_jsonl_manifest(&leaf, 1, 5)?; + crate::store::catalog::manifest::write_compact_jsonl_manifest(&nested, 1, 99)?; let snapshot = DatasetCatalogSnapshot::discover( vec![DatasetMount::default(leaf.to_string_lossy())?], @@ -407,14 +407,14 @@ async fn updating_one_leaf_manifest_does_not_require_rewriting_parent_branch() - let warehouse = temp.path().join("warehouse"); let leaf = warehouse.join("codex_jsonl"); fs::create_dir_all(&leaf)?; - crate::store::chronicle_manifest::atomic_write_manifest( + crate::store::catalog::manifest::atomic_write_manifest( &warehouse, &crate::store::ChronicleManifest::branch(), )?; - crate::store::chronicle_manifest::write_compact_jsonl_manifest(&leaf, 1, 10)?; + crate::store::catalog::manifest::write_compact_jsonl_manifest(&leaf, 1, 10)?; let parent_before = fs::read(warehouse.join("chronicle.manifest"))?; - crate::store::chronicle_manifest::write_compact_jsonl_manifest(&leaf, 2, 42)?; + crate::store::catalog::manifest::write_compact_jsonl_manifest(&leaf, 2, 42)?; let parent_after = fs::read(warehouse.join("chronicle.manifest"))?; assert_eq!( parent_before, parent_after, diff --git a/crates/persisting-pchronicle/src/store/compact_jsonl.rs b/crates/persisting-pchronicle/src/store/compact_jsonl.rs index f9d81cb18..d0ab8f109 100644 --- a/crates/persisting-pchronicle/src/store/compact_jsonl.rs +++ b/crates/persisting-pchronicle/src/store/compact_jsonl.rs @@ -179,12 +179,8 @@ impl CompactJsonlStore { validate_dataset_schema(&dataset)?; let version = dataset.version_id(); let record_count = dataset.count_rows(None).await? as u64; - crate::store::chronicle_manifest::write_compact_jsonl_manifest( - root, - version, - record_count, - )?; - crate::store::chronicle_manifest::load_manifest(root)? + crate::store::catalog::manifest::write_compact_jsonl_manifest(root, version, record_count)?; + crate::store::catalog::manifest::load_manifest(root)? .context("chronicle.manifest missing after publish") } @@ -200,18 +196,18 @@ impl CompactJsonlStore { return Ok(None); } let version = dataset.version_id(); - if let Some(manifest) = crate::store::chronicle_manifest::try_load_manifest(root) - && crate::store::chronicle_manifest::compact_jsonl_manifest_matches(&manifest, version) + if let Some(manifest) = crate::store::catalog::manifest::try_load_manifest(root) + && crate::store::catalog::manifest::compact_jsonl_manifest_matches(&manifest, version) { return Ok(Some(manifest)); } let record_count = dataset.count_rows(None).await? as u64; - match crate::store::chronicle_manifest::write_compact_jsonl_manifest( + match crate::store::catalog::manifest::write_compact_jsonl_manifest( root, version, record_count, ) { - Ok(()) => Ok(crate::store::chronicle_manifest::try_load_manifest(root)), + Ok(()) => Ok(crate::store::catalog::manifest::try_load_manifest(root)), Err(error) => { tracing::warn!( target: "persisting_pchronicle::compact_jsonl", @@ -1016,7 +1012,7 @@ mod tests { b"{\"id\":\"a\",\"timestamp\":1}\n{\"id\":\"b\",\"timestamp\":2}\n", )?; CompactJsonlStore::import_path(&input, &dataset, &CompactJsonlOptions::default()).await?; - fs::remove_file(crate::store::chronicle_manifest::manifest_path(&dataset))?; + fs::remove_file(crate::store::catalog::manifest::manifest_path(&dataset))?; assert!(crate::store::load_manifest(&dataset)?.is_none()); let ensured = CompactJsonlStore::ensure_manifest(&dataset) diff --git a/crates/persisting-pchronicle/src/store/mod.rs b/crates/persisting-pchronicle/src/store/mod.rs index 0e2a657e5..e9939dda8 100644 --- a/crates/persisting-pchronicle/src/store/mod.rs +++ b/crates/persisting-pchronicle/src/store/mod.rs @@ -16,7 +16,6 @@ mod cas_store; #[cfg(feature = "lance-store")] mod catalog; #[cfg(feature = "lance-store")] -mod chronicle_manifest; #[cfg(feature = "lance-store")] mod compact_jsonl; #[cfg(feature = "lance-store")] @@ -42,14 +41,17 @@ mod inspect; #[cfg(feature = "lance-store")] mod local_query_manifest; #[cfg(feature = "lance-store")] -mod location; #[cfg(feature = "lance-store")] pub(crate) mod object_store_io_gate; #[cfg(feature = "lance-store")] pub(crate) mod opendal_store; #[cfg(feature = "lance-store")] +pub mod persistent_cache; +#[cfg(feature = "lance-store")] mod query_engine; #[cfg(feature = "lance-store")] +pub use persistent_cache::PersistentCache; +#[cfg(feature = "lance-store")] mod root_write_lock; #[cfg(feature = "lance-store")] mod run_control; @@ -68,23 +70,30 @@ pub use attempt_registry::{AttemptRecord, AttemptRecordState, AttemptRegistry}; #[cfg(feature = "lance-store")] pub use cas_store::unix_now_ms as attempt_registry_now_ms; #[cfg(feature = "lance-store")] -pub use catalog::{ - CATALOG_SOURCES_TABLE, CATALOG_TRAJECTORIES_TABLE, CatalogDataset, CatalogErrorPolicy, - CatalogEventProvenance, CatalogEventView, CatalogNamespace, CatalogPage, - CatalogProjectionStatus, CatalogSnapshotOptions, CatalogSourceDescription, CatalogSourceKind, - CatalogSourceRevision, CatalogSourceStatus, CatalogStorylineKey, CatalogTrajectoryBundle, - DEFAULT_DATASET_NAME, DEFAULT_MAX_EVENT_FALLBACK_BYTES, DEFAULT_MAX_EVENT_FALLBACK_ROWS, - DatasetCatalogSnapshot, DatasetMount, DiscoveredSource, NamespacePath, QueryScope, +pub use catalog::location::{ + DatasetLocation, DatasetLocationKind, ImportableObjectEvent, PathListEntry, PathListKind, + ShallowNavEntry, }; #[cfg(feature = "lance-store")] #[allow(unused_imports)] -pub use chronicle_manifest::{ +pub use catalog::manifest::{ CHRONICLE_MANIFEST_FILE, ChronicleManifest, ManifestKind, ManifestStats, STORYLINE_FORMAT, atomic_write_manifest, compact_jsonl_manifest_matches, load_manifest, load_manifest_at_uri, try_load_manifest, write_compact_jsonl_manifest, write_storyline_manifest, write_storyline_manifest_at_uri, }; #[cfg(feature = "lance-store")] +pub use catalog::{ + CATALOG_SOURCES_TABLE, CATALOG_TRAJECTORIES_TABLE, CachedDataset, CatalogDataset, + CatalogErrorPolicy, CatalogEventProvenance, CatalogEventView, CatalogNamespace, CatalogPage, + CatalogProjectionStatus, CatalogSnapshotOptions, CatalogSourceDescription, CatalogSourceKind, + CatalogSourceRevision, CatalogSourceStatus, CatalogStorylineKey, CatalogTrajectoryBundle, + DEFAULT_DATASET_NAME, DEFAULT_MAX_EVENT_FALLBACK_BYTES, DEFAULT_MAX_EVENT_FALLBACK_ROWS, + Dataset, DatasetCatalogSnapshot, DatasetMount, DatasetResolver, DiscoveredSource, + LocationSummary, ManifestCache, ManifestListing, ManifestReadMode, NamespacePath, QueryScope, + ResolveMode, ResolveTarget, +}; +#[cfg(feature = "lance-store")] pub use compact_jsonl::{ CompactJsonlBuildPhase, CompactJsonlColumn, CompactJsonlImportEvent, CompactJsonlOffload, CompactJsonlOptions, CompactJsonlRecord, CompactJsonlStore, @@ -125,11 +134,6 @@ pub(crate) use local_query_manifest::{ LocalQueryInputFile, LocalQueryManifest, LocalQueryManifestOptions, }; #[cfg(feature = "lance-store")] -pub use location::{ - DatasetLocation, DatasetLocationKind, ImportableObjectEvent, PathListEntry, PathListKind, - ShallowNavEntry, -}; -#[cfg(feature = "lance-store")] pub use query_engine::{ ChronicleQueryEngine, ChronicleQueryExecutionOptions, DEFAULT_QUERY_MEMORY_LIMIT_BYTES, ExternalTableFormat, ExternalTableSpec, IntrospectedField, IntrospectedTable, diff --git a/crates/persisting-pchronicle/src/store/persistent_cache.rs b/crates/persisting-pchronicle/src/store/persistent_cache.rs new file mode 100644 index 000000000..e2a43ae38 --- /dev/null +++ b/crates/persisting-pchronicle/src/store/persistent_cache.rs @@ -0,0 +1,154 @@ +//! Rebuildable Lance-backed key/value cache. It is never authoritative. +use anyhow::{Context, Result}; +use futures::TryStreamExt; +use lance::Dataset; +use lance::dataset::{InsertBuilder, MergeInsertBuilder, WhenMatched, WhenNotMatched}; +use lance::deps::arrow_array::{Array, RecordBatch, RecordBatchIterator, StringArray}; +use lance::deps::arrow_schema::{DataType, Field, Schema}; +use serde::{Serialize, de::DeserializeOwned}; +use std::collections::HashMap; +use std::path::{Path, PathBuf}; +use std::sync::Arc; +use tokio::sync::RwLock; +const SCHEMA_VERSION: &str = "persistent-cache-v1"; +pub struct PersistentCache { + path: PathBuf, + values: RwLock>, + disk_lock: Option, +} +impl PersistentCache +where + K: Eq + std::hash::Hash + Serialize + DeserializeOwned, + V: Serialize + DeserializeOwned, +{ + pub async fn open(path: PathBuf) -> Self { + let lock = (|| -> Result { + std::fs::create_dir_all(path.parent().context("cache parent")?)?; + let file = std::fs::OpenOptions::new() + .create(true) + .truncate(false) + .read(true) + .write(true) + .open(path.with_extension("lock"))?; + fs2::FileExt::try_lock_exclusive(&file)?; + Ok(file) + })(); + let mut cache = Self { + path, + values: RwLock::new(HashMap::new()), + disk_lock: lock.ok(), + }; + if cache.disk_lock.is_some() { + match cache.load().await { + Ok(values) => *cache.values.get_mut() = values, + Err(error) => { + tracing::warn!(target: "pchronicle.serve", error = %error, "persistent cache unreadable; rebuilding"); + let _ = if cache.path.is_dir() { + std::fs::remove_dir_all(&cache.path) + } else { + std::fs::remove_file(&cache.path) + }; + } + } + } + cache + } + pub fn path(&self) -> &Path { + &self.path + } + pub fn writable(&self) -> bool { + self.disk_lock.is_some() + } + pub async fn values(&self) -> HashMap + where + K: Clone, + V: Clone, + { + self.values.read().await.clone() + } + pub async fn upsert(&self, key: &K, value: &V, removed: &[K]) -> Result<()> { + if !self.writable() { + return Ok(()); + } + let schema = Arc::new(Schema::new(vec![ + Field::new("key", DataType::Utf8, false), + Field::new("payload", DataType::Utf8, false), + Field::new("schema_version", DataType::Utf8, false), + ])); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(StringArray::from(vec![serde_json::to_string(key)?])) as _, + Arc::new(StringArray::from(vec![serde_json::to_string(value)?])) as _, + Arc::new(StringArray::from(vec![SCHEMA_VERSION])) as _, + ], + )?; + if self.path.exists() { + let mut dataset = Dataset::open(self.path.to_string_lossy().as_ref()).await?; + for chunk in removed.chunks(128) { + let keys = chunk + .iter() + .map(|key| { + Ok(format!( + "'{}'", + serde_json::to_string(key)?.replace("'", "''") + )) + }) + .collect::>>()? + .join(","); + dataset.delete(&format!("key IN ({keys})")).await?; + } + MergeInsertBuilder::try_new(Arc::new(dataset), vec!["key".into()])? + .when_matched(WhenMatched::UpdateAll) + .when_not_matched(WhenNotMatched::InsertAll) + .try_build()? + .execute_reader(Box::new(RecordBatchIterator::new(vec![Ok(batch)], schema))) + .await?; + } else { + InsertBuilder::new(self.path.to_string_lossy().as_ref()) + .execute(vec![batch]) + .await?; + } + Ok(()) + } + async fn load(&self) -> Result> { + if !self.path.exists() { + return Ok(HashMap::new()); + } + let dataset = Dataset::open(self.path.to_string_lossy().as_ref()).await?; + let batches: Vec = dataset + .scan() + .try_into_stream() + .await? + .try_collect() + .await?; + let mut values = HashMap::new(); + for batch in batches { + let keys = text_column(&batch, "key")?; + let payloads = text_column(&batch, "payload")?; + let versions = text_column(&batch, "schema_version")?; + for row in 0..batch.num_rows() { + anyhow::ensure!( + versions[row] == SCHEMA_VERSION, + "unsupported persistent cache schema" + ); + values.insert( + serde_json::from_str(&keys[row])?, + serde_json::from_str(&payloads[row])?, + ); + } + } + Ok(values) + } +} +fn text_column(batch: &RecordBatch, name: &str) -> Result> { + let array = batch + .column(batch.schema().index_of(name)?) + .as_any() + .downcast_ref::() + .with_context(|| format!("persistent cache column {name} must be Utf8"))?; + anyhow::ensure!(array.null_count() == 0, "null persistent cache field"); + Ok((0..array.len()) + .map(|index| array.value(index).to_owned()) + .collect()) +} diff --git a/docs/src/en/pchronicle/guides/serve.md b/docs/src/en/pchronicle/guides/serve.md index d55de5d47..b3ab9ff15 100644 --- a/docs/src/en/pchronicle/guides/serve.md +++ b/docs/src/en/pchronicle/guides/serve.md @@ -126,6 +126,20 @@ and arbitrary filesystem access are not exposed through the API. Refreshes replace the readable view only after the replacement is ready; a failed refresh keeps the previous view available. +Runs lists allow a short freshness delay. Successful summaries are reused by +Dataset and file scope. After 30 seconds, the next visit returns cached results +and starts a background refresh. At most one automatic Runs refresh runs per +process. Failures retain the previous result and wait at least 30 seconds before +retrying. The summary cache retains at most 32 scopes and 64 MiB in memory; +first access, process restart, or eviction still requires source discovery. +The existing Datasets browse cache remains persisted in local Lance storage. +Unscoped Runs requests also reuse their Catalog and refresh it in the background. +**Refresh** explicitly rebuilds the Catalog and invalidates scoped summaries on +success; an older in-flight build cannot republish an invalidated cache entry. +CLI `query`, query workers, and Gateway live reads bypass the summary cache. + +During startup the Web UI uses the lightweight `query/tables?ui=1` Catalog view. Dataset names and source counts come from the local Browse/manifest cache without triggering accurate Catalog discovery. It returns the `ui-cache` marker; Analysis compilation and SQL execution still bind to the current accurate Snapshot on the server. CLI and API requests without `ui=1` do not use this view. + ## Logs and failed requests `pchronicle serve` writes Warehouse request logs to stderr at `--log-level` diff --git a/docs/src/zh/pchronicle/guides/serve.md b/docs/src/zh/pchronicle/guides/serve.md index b6a706e9f..a04de435a 100644 --- a/docs/src/zh/pchronicle/guides/serve.md +++ b/docs/src/zh/pchronicle/guides/serve.md @@ -114,6 +114,16 @@ readiness 记录;Control 凭据不会写入 stderr。 挂载的 Dataset 和 HTTP 操作均为只读。API 不暴露 import、export、maintenance 或任意文件 访问。刷新只会在新视图准备完成后替换当前可读视图;刷新失败时,旧视图继续可用。 +Runs 列表允许短暂滞后:按 Dataset 和文件范围复用最近成功的摘要,30 秒后再次访问时先 +返回旧结果,再触发后台刷新。全进程最多一个自动 Runs 刷新任务;失败保留旧结果,并至少 +等待 30 秒再重试。摘要缓存最多保留 32 个范围、64 MiB;它只驻留内存,首次访问、进程重启 +或缓存淘汰后仍需发现 source。已有的 Datasets 目录浏览缓存仍使用本地 Lance 持久化。 +未限定 Dataset 的 Runs 请求同样复用已有 Catalog,过期时后台更新。点击 **Refresh** 会 +强制刷新 Catalog,成功后使旧的分范围摘要失效;刷新前的任务不能重新发布旧缓存。 +命令行 `query`、查询 worker 和 Gateway 实时读取不使用这层摘要缓存。 + +Web UI 启动阶段使用 `query/tables?ui=1` 的轻量 Catalog 视图:Dataset 名称和 source 统计优先来自本地 Browse/manifest cache,不触发准确 Catalog discovery。它返回 `ui-cache` 标识;真正的 Analysis 编译和 SQL 执行仍在服务端绑定当前准确 Snapshot。CLI 和未带 `ui=1` 的 API 不使用此视图。 + ## 日志与失败请求 `pchronicle serve` 把 Warehouse 请求日志写到 stderr,级别由 `--log-level` 控制(默认 diff --git a/pchronicle-web/assets/span-timeline.css b/pchronicle-web/assets/span-timeline.css index af69c268d..aeeb82552 100644 --- a/pchronicle-web/assets/span-timeline.css +++ b/pchronicle-web/assets/span-timeline.css @@ -96,6 +96,7 @@ .trace-root, .span-row { + position: relative; border: 0; } diff --git a/pchronicle-web/src/api.rs b/pchronicle-web/src/api.rs index 5692a1662..521209b4f 100644 --- a/pchronicle-web/src/api.rs +++ b/pchronicle-web/src/api.rs @@ -267,7 +267,7 @@ pub async fn compile_analysis( pub async fn query_catalog() -> Result { json_checked( - with_catalog_headers(Request::get("/api/query/tables")) + with_catalog_headers(Request::get("/api/query/tables?ui=true")) .send() .await, ) diff --git a/pchronicle-web/src/components.rs b/pchronicle-web/src/components.rs index d9bfb976f..aa8b04a63 100644 --- a/pchronicle-web/src/components.rs +++ b/pchronicle-web/src/components.rs @@ -967,10 +967,10 @@ fn CompactSpanRow( } OccupancyTrack { bars, expose_range, focus_left, caption: caption.clone(), title: "{caption} · {meta}", exposed_ids, expanded_turn_id, hovered_ids } div { class: "span-evidence-count", if event_refs > 0 { span { class: "span-count-chip event", "{event_refs} events" } } } - if !embedded { - if let Some(id) = drawer_id { - button { class: "pc2-conversation-drawer-button", title: "Open conversation as AgenticMD", aria_label: "Open {drawer_label} as AgenticMD", onclick: move |event| { event.prevent_default(); event.stop_propagation(); on_open_drawer.call((id, drawer_label.clone(), drawer_ids.clone())); }, "↗" } - } + } + if !embedded { + if let Some(id) = drawer_id { + button { class: "pc2-conversation-drawer-button", title: "Open conversation as AgenticMD", aria_label: "Open {drawer_label} as AgenticMD", onclick: move |event| { event.prevent_default(); event.stop_propagation(); on_open_drawer.call((id, drawer_label.clone(), drawer_ids.clone())); }, "↗" } } } if row_open {