Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
47ca681
feat(catalog): introduce catalog state management and consistency tra…
reiase Sep 14, 2026
2791a0c
feat(pchronicle): enhance caching and storage mechanisms
reiase Sep 15, 2026
e6bb07c
feat(catalog): enhance manifest caching and dataset browsing capabili…
reiase Sep 15, 2026
c7ecd70
fix(workspace): improve error handling in dataset loading
reiase Sep 15, 2026
af441ce
feat(server): implement timeout for run scans and enhance error handling
reiase Sep 15, 2026
575a4d2
feat(server): implement structured timeout handling for run operations
reiase Sep 15, 2026
0a2c855
feat(workspace): add run generation tracking to load_runs function
reiase Sep 15, 2026
6e2f1ab
feat(catalog): enhance dataset browsing and manifest refresh capabili…
reiase Sep 15, 2026
07fb8da
feat(object_store): enhance concurrency control and error handling
reiase Sep 15, 2026
43a9502
feat(model): enhance dataset browsing and status management
reiase Sep 15, 2026
5fab762
fix(workspace): prevent request loop in load_runs function
reiase Sep 15, 2026
c87c983
feat(store): enhance object store configuration and background I/O ma…
reiase Sep 16, 2026
04c4ca8
feat(request_progress): implement request diagnostics and progress tr…
reiase Sep 16, 2026
f1d7b05
fix(workspace): optimize run loading behavior in detail view
reiase Sep 16, 2026
06b2012
feat(worker_pool): enhance worker management and idle cleanup
reiase Sep 16, 2026
2007cfc
Fix search previews on run pages
reiase Sep 16, 2026
8d6d1de
Parallelize remote object store reads
reiase Sep 16, 2026
2e1b738
Enhance object store management and retry mechanisms
reiase Sep 18, 2026
cef1121
Refactor object store and cache management for improved readability
reiase Sep 18, 2026
6a8abd1
Refactor error handling and improve cache management in pChronicle
reiase Sep 18, 2026
109e011
Enhance object store configuration and cache management
reiase Sep 18, 2026
d046ea2
Merge branch 'main' into feature/better_error
reiase Sep 18, 2026
3a8144b
Refactor catalog access control and dataset handling
reiase Sep 18, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@ lance-core = "11.0.0"
lance-file = "11.0.0"
lance-datafusion = "11.0.0"
lance-index = { version = "11.0.0", features = ["tokenizer-jieba"] }
lance-io = "11.0.0"
lance-linalg = "11.0.0"
lance-table = "11.0.0"
libc = "0.2"
Expand All @@ -75,6 +76,7 @@ opendal = { version = "0.57.0", default-features = true, features = [
"services-s3",
"services-tos",
] }
object_store = "0.13.2"
persisting-agentctl = { path = "crates/persisting-agentctl" }
persisting-events = { path = "crates/persisting-events" }
persisting-gateway = { path = "crates/persisting-gateway" }
Expand Down
10 changes: 8 additions & 2 deletions crates/persisting-pchronicle-cli/src/gateway_partition.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
//! Safe, bounded physical partitioning below one logical Gateway Dataset.

use std::collections::HashMap;
use std::sync::Mutex;
use std::sync::{Mutex, MutexGuard};

use anyhow::{Context, Result};
use chrono::{DateTime, Datelike, Timelike, Utc};
Expand All @@ -10,6 +10,12 @@ const MAX_TEMPLATE_BYTES: usize = 256;
const MAX_TEMPLATE_SEGMENTS: usize = 16;
const MAX_USER_SEGMENT_BYTES: usize = 80;

fn lock_recover<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
mutex
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}

#[derive(Debug, Clone, PartialEq, Eq)]
enum Segment {
Literal(String),
Expand Down Expand Up @@ -149,7 +155,7 @@ impl GatewayPartitionRouter {
let Some(split) = &self.split else {
return self.dataset_uri.clone();
};
let mut routes = self.routes.lock().unwrap();
let mut routes = lock_recover(&self.routes);
routes
.entry(route_key.to_string())
.or_insert_with(|| {
Expand Down
94 changes: 60 additions & 34 deletions crates/persisting-pchronicle-cli/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -121,7 +121,9 @@ impl Cli {
/// invocations return None and retain their existing execution path.
pub fn run_catalog_worker_before_runtime(cli: &Cli) -> Option<Result<()>> {
match &cli.command {
Command::Serve(args) if args.catalog_query_worker => Some(server::catalog_worker::run()),
Command::Serve(args) if args.catalog_query_worker => {
Some(server::catalog_worker::run(cli.log_level))
}
_ => None,
}
}
Expand Down Expand Up @@ -172,6 +174,18 @@ pub enum LogLevel {
Debug,
}

impl LogLevel {
/// The `--log-level` value that parses back to this variant.
pub(crate) fn as_arg(self) -> &'static str {
match self {
Self::Error => "error",
Self::Warn => "warn",
Self::Info => "info",
Self::Debug => "debug",
}
}
}

struct DiagnosticWriter<'a> {
level: LogLevel,
inner: &'a mut dyn Write,
Expand Down Expand Up @@ -3589,51 +3603,63 @@ pub(crate) async fn find_expression_predicate_for_dataset(
for predicate in text_predicates {
let mut matches = Vec::<String>::new();
let mut searched = false;
let mut jobs = Vec::new();
for dataset in snapshot.datasets() {
if dataset_filter.is_some_and(|filter| dataset.mount.name != filter) {
continue;
}
for source in &dataset.sources {
if source.status != CatalogSourceStatus::Ready
|| source.kind != CatalogSourceKind::Store
|| source_filter.is_some_and(|filter| filter != source.file)
|| source_filter.is_some_and(|filter| {
filter != source.file && !source.file.starts_with(&format!("{filter}/"))
})
{
continue;
}
let Some(paths) =
if let Some(paths) =
snapshot.storyline_table_paths(&dataset.mount.name, &source.file)?
else {
continue;
};
match search_storyline_step_matches_fts_in_columns(
&paths,
&predicate.query,
predicate.field.columns(),
)
.await
{
Ok(step_matches) => {
available = true;
searched = true;
let source_predicate = predicate
.field
.source_predicate()
.map(|value| format!(" AND ({value})"))
.unwrap_or_default();
matches.extend(step_matches.into_iter().map(|(document_id, step_id)| {
format!(
"(_file_ = {} AND document_id = {} AND step_id = {}{})",
sql_string(&source.file),
sql_string(&document_id),
step_id,
source_predicate,
)
}));
}
Err(error) => errors.push(format!(
"FTS unavailable for {} / {}: {error:#}",
dataset.mount.name, source.file
)),
jobs.push((dataset.mount.name.clone(), source.file.clone(), paths));
}
}
}
let query = predicate.query.clone();
let columns = predicate.field.columns();
let results = stream::iter(jobs)
.map(|(dataset, file, paths)| {
let query = query.clone();
async move {
let result =
search_storyline_step_matches_fts_in_columns(&paths, &query, columns).await;
(dataset, file, result)
}
})
.buffer_unordered(4)
.collect::<Vec<_>>()
.await;
for (dataset, file, result) in results {
match result {
Ok(step_matches) => {
available = true;
searched = true;
let source_predicate = predicate
.field
.source_predicate()
.map(|value| format!(" AND ({value})"))
.unwrap_or_default();
matches.extend(step_matches.into_iter().map(|(document_id, step_id)| {
format!(
"(_file_ = {} AND document_id = {} AND step_id = {}{})",
sql_string(&file),
sql_string(&document_id),
step_id,
source_predicate,
)
}));
}
Err(error) => {
errors.push(format!("FTS unavailable for {dataset} / {file}: {error:#}"))
}
}
}
Expand Down
Loading
Loading