From f0ff9938973483e157fd3ad588efaab2441c86eb Mon Sep 17 00:00:00 2001 From: Reiase Date: Sun, 13 Sep 2026 06:48:03 +0800 Subject: [PATCH 1/3] Refactor dataset location handling and enhance caching mechanisms - Updated `DatasetLocation` methods to improve dataset preview and navigation functionality, including the introduction of a new `list_nav_children` method. - Enhanced the `list_dataset_preview` method to accept a dataset kind parameter, streamlining the preview retrieval process. - Introduced `discover_cached_candidates` and `discover_cached_candidate_at` functions to optimize candidate discovery from cached sources. - Updated `CatalogSnapshot` methods to support cached file discovery, improving performance during catalog operations. - Enhanced request logging with metrics tracking for better performance insights in the CLI server. - Improved UI caching mechanisms to retain successful summaries and manage background refreshes effectively. This commit aims to enhance the overall efficiency and user experience of dataset handling and querying in the application. --- .../src/server/acceleration.rs | 17 + .../src/server/mod.rs | 561 +++++++++++++----- .../src/server/query_admission.rs | 297 +++++++++- .../src/server/request_log.rs | 50 +- .../src/server/tests.rs | 131 ++++ .../src/server/ui_cache.rs | 28 +- .../src/store/catalog/discovery.rs | 68 +++ .../src/store/catalog/mod.rs | 36 +- .../src/store/location.rs | 154 +++-- docs/src/en/pchronicle/guides/serve.md | 14 + docs/src/zh/pchronicle/guides/serve.md | 10 + pchronicle-web/assets/span-timeline.css | 1 + pchronicle-web/src/api.rs | 2 +- pchronicle-web/src/components.rs | 8 +- 14 files changed, 1146 insertions(+), 231 deletions(-) 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..0759b893f 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,26 @@ 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( + let snapshot = Arc::new(if cached_files.is_empty() { DatasetCatalogSnapshot::discover_scoped( config.datasets.clone(), config.default_dataset.clone(), config.catalog_options, scope, ) - .await?, - ); + .await? + } else { + DatasetCatalogSnapshot::discover_scoped_from_cached_files( + config.datasets.clone(), + config.default_dataset.clone(), + config.catalog_options, + scope, + cached_files, + ) + .await? + }); let engine = Arc::new( snapshot .clone() @@ -547,16 +560,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 +612,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 +643,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 +668,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 +681,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 +689,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 +1099,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 +1115,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 +1356,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 +1367,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 +1385,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 +1399,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 +1510,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 +1524,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 +1710,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 +1773,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 +1912,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 +1935,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 +1946,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 +2028,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 +2055,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 +2080,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 +2127,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 +2144,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 +2156,7 @@ async fn explorer_turns( } else { false }; + metrics.record("turn_fts_probe", phase); let mut search_mode = if query .q .as_deref() @@ -2010,6 +2182,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 +2192,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 +2248,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 +2261,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 +2280,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 +2294,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 +2549,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 +2751,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 +2915,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..c4593f377 100644 --- a/crates/persisting-pchronicle-cli/src/server/ui_cache.rs +++ b/crates/persisting-pchronicle-cli/src/server/ui_cache.rs @@ -166,6 +166,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; @@ -383,6 +408,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 +417,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) } }; diff --git a/crates/persisting-pchronicle/src/store/catalog/discovery.rs b/crates/persisting-pchronicle/src/store/catalog/discovery.rs index 239cd3b78..39a81262f 100644 --- a/crates/persisting-pchronicle/src/store/catalog/discovery.rs +++ b/crates/persisting-pchronicle/src/store/catalog/discovery.rs @@ -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/catalog/mod.rs b/crates/persisting-pchronicle/src/store/catalog/mod.rs index 7dccb75c3..fbb0fab37 100644 --- a/crates/persisting-pchronicle/src/store/catalog/mod.rs +++ b/crates/persisting-pchronicle/src/store/catalog/mod.rs @@ -17,8 +17,8 @@ use provider::*; 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 +340,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 +349,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 +377,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 +414,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/location.rs b/crates/persisting-pchronicle/src/store/location.rs index ae544c1bf..fc7efa9cb 100644 --- a/crates/persisting-pchronicle/src/store/location.rs +++ b/crates/persisting-pchronicle/src/store/location.rs @@ -6,6 +6,7 @@ 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; @@ -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(), @@ -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, @@ -1231,43 +1233,75 @@ mod tests { crate::store::chronicle_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/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 { From f53e58c887d7f9115b74e076aeb3e0f1564c99a5 Mon Sep 17 00:00:00 2001 From: Reiase Date: Sun, 13 Sep 2026 09:54:30 +0800 Subject: [PATCH 2/3] Refactor storage module and enhance catalog management - Updated the `storage.rs` file to replace references from `chronicle_manifest` to `catalog::manifest`, improving modularity and clarity. - Introduced a new `PersistentCache` module to manage key/value caching, enhancing performance and data retrieval efficiency. - Added `DatasetLocation` and `DatasetResolver` structures to streamline dataset handling and improve location management. - Enhanced the `CompactJsonlStore` to utilize the new manifest handling methods, ensuring consistency across data operations. - Introduced new methods for managing dataset locations and improved error handling in manifest loading. This commit aims to improve the overall architecture and performance of the data storage and catalog management systems. --- .../src/server/mod.rs | 43 +-- .../src/server/ui_cache.rs | 231 +++++--------- crates/persisting-pchronicle/src/storage.rs | 48 +-- .../src/store/catalog/discovery.rs | 6 +- .../src/store/{ => catalog}/location.rs | 12 +- .../manifest.rs} | 4 +- .../src/store/catalog/manifest_cache.rs | 286 ++++++++++++++++++ .../src/store/catalog/mod.rs | 19 ++ .../src/store/catalog/resolver.rs | 221 ++++++++++++++ .../src/store/catalog/tests.rs | 26 +- .../src/store/compact_jsonl.rs | 18 +- crates/persisting-pchronicle/src/store/mod.rs | 33 +- .../src/store/persistent_cache.rs | 154 ++++++++++ 13 files changed, 853 insertions(+), 248 deletions(-) rename crates/persisting-pchronicle/src/store/{ => catalog}/location.rs (99%) rename crates/persisting-pchronicle/src/store/{chronicle_manifest.rs => catalog/manifest.rs} (98%) create mode 100644 crates/persisting-pchronicle/src/store/catalog/manifest_cache.rs create mode 100644 crates/persisting-pchronicle/src/store/catalog/resolver.rs create mode 100644 crates/persisting-pchronicle/src/store/persistent_cache.rs diff --git a/crates/persisting-pchronicle-cli/src/server/mod.rs b/crates/persisting-pchronicle-cli/src/server/mod.rs index 0759b893f..3ce0dc794 100644 --- a/crates/persisting-pchronicle-cli/src/server/mod.rs +++ b/crates/persisting-pchronicle-cli/src/server/mod.rs @@ -502,24 +502,33 @@ async fn build_scoped_query_runtime( scope: persisting_pchronicle::storage::QueryScope, cached_files: Vec, ) -> anyhow::Result> { - let snapshot = Arc::new(if cached_files.is_empty() { - 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 { - DatasetCatalogSnapshot::discover_scoped_from_cached_files( - config.datasets.clone(), - config.default_dataset.clone(), - config.catalog_options, - scope, - cached_files, - ) - .await? - }); + 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() diff --git a/crates/persisting-pchronicle-cli/src/server/ui_cache.rs b/crates/persisting-pchronicle-cli/src/server/ui_cache.rs index c4593f377..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, @@ -321,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, @@ -330,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! { @@ -344,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() @@ -376,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. @@ -387,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); } } } @@ -437,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) { @@ -528,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() @@ -636,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)); } @@ -661,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 @@ -682,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 { @@ -732,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; @@ -741,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)); } @@ -767,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()) @@ -790,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 39a81262f..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); diff --git a/crates/persisting-pchronicle/src/store/location.rs b/crates/persisting-pchronicle/src/store/catalog/location.rs similarity index 99% rename from crates/persisting-pchronicle/src/store/location.rs rename to crates/persisting-pchronicle/src/store/catalog/location.rs index fc7efa9cb..e0e929ff7 100644 --- a/crates/persisting-pchronicle/src/store/location.rs +++ b/crates/persisting-pchronicle/src/store/catalog/location.rs @@ -9,7 +9,7 @@ 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)] @@ -35,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, @@ -43,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, @@ -282,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")); } @@ -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 { @@ -1230,7 +1230,7 @@ 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 remote = DatasetLocation::parse(&format!( 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 fbb0fab37..61f101834 100644 --- a/crates/persisting-pchronicle/src/store/catalog/mod.rs +++ b/crates/persisting-pchronicle/src/store/catalog/mod.rs @@ -7,13 +7,32 @@ 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::{ 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..514e7d1d7 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,13 +41,15 @@ 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; +pub use persistent_cache::PersistentCache; #[cfg(feature = "lance-store")] mod root_write_lock; #[cfg(feature = "lance-store")] @@ -68,23 +69,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 +133,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()) +} From ad50888ef09c4c148f017bc874671d5c9c2dc56e Mon Sep 17 00:00:00 2001 From: Reiase Date: Sun, 13 Sep 2026 10:08:04 +0800 Subject: [PATCH 3/3] Expose `PersistentCache` in the `lance-store` feature module to enhance caching capabilities. This change allows for better management of key/value caching within the storage module, improving performance and data retrieval efficiency. --- crates/persisting-pchronicle/src/store/mod.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/crates/persisting-pchronicle/src/store/mod.rs b/crates/persisting-pchronicle/src/store/mod.rs index 514e7d1d7..e9939dda8 100644 --- a/crates/persisting-pchronicle/src/store/mod.rs +++ b/crates/persisting-pchronicle/src/store/mod.rs @@ -49,6 +49,7 @@ pub(crate) mod opendal_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;