Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
50 changes: 42 additions & 8 deletions crates/persisting-pchronicle-cli/src/exchange/import.rs
Original file line number Diff line number Diff line change
Expand Up @@ -587,6 +587,12 @@ pub(crate) async fn run_import(
progress.notice(&line)?;
}
progress.flush_log(stderr)?;
if let Some(wal) = &wal
&& let Ok(guard) = wal.lock()
&& let Err(error) = guard.remove()
{
tracing::warn!(target: "pchronicle.import", %error, "failed to remove completed import WAL");
}
Ok(())
}

Expand Down Expand Up @@ -1938,9 +1944,20 @@ pub(crate) async fn finalize_storyline_import_indexes(
store: &StorylineLanceStore,
progress: &mut CliProgress,
) -> Result<()> {
progress
.stage(StageId::Commit)
.set_current("optimize indices (final)");
let commit = progress.stage(StageId::Commit);
commit.set_current("data committed; building indices");
let heartbeat = commit.clone();
let started = std::time::Instant::now();
let heartbeat_task = tokio::spawn(async move {
let mut tick = tokio::time::interval(std::time::Duration::from_secs(2));
loop {
tick.tick().await;
heartbeat.set_activity_override(&format!(
"building indices · elapsed={}s",
started.elapsed().as_secs()
));
}
});
let _index_progress = progress.attach_index_progress();
store
.maintain(&persisting_pchronicle::storage::LanceMaintenanceOptions {
Expand All @@ -1951,9 +1968,9 @@ pub(crate) async fn finalize_storyline_import_indexes(
})
.await
.context("finalize Storyline indexes after progressive import")?;
progress
.stage(StageId::Commit)
.set_current("optimize indices done");
heartbeat_task.abort();
commit.set_activity_override("indices built");
commit.set_current("indices built");
Ok(())
}

Expand All @@ -1968,7 +1985,20 @@ pub(crate) async fn commit_storyline_import_batch(
anyhow::ensure!(!batch.is_empty(), "storyline import commit batch is empty");
let batch_len = batch.len() as u64;
commit.set_current(format!("batch={batch_len}"));
let report = match append_generation.as_deref() {
let heartbeat = commit.clone();
let heartbeat_started = std::time::Instant::now();
let heartbeat_task = tokio::spawn(async move {
let mut tick = tokio::time::interval(std::time::Duration::from_secs(2));
loop {
tick.tick().await;
heartbeat.set_activity_override(&format!(
"writing Lance snapshot · batch={batch_len} · elapsed={}s",
heartbeat_started.elapsed().as_secs()
));
}
});
let report_result: Result<_> = async {
Ok(match append_generation.as_deref() {
Some(generation) => {
tracing::info!(
committed_before = committed_storylines,
Expand Down Expand Up @@ -2010,7 +2040,11 @@ pub(crate) async fn commit_storyline_import_batch(
)
})?
}
};
})
}
.await;
heartbeat_task.abort();
let report = report_result?;
anyhow::ensure!(
report.storylines as u64 == batch_len,
"storyline import batch report does not match batch size"
Expand Down
12 changes: 12 additions & 0 deletions crates/persisting-pchronicle-cli/src/exchange/wal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -160,6 +160,15 @@ impl ImportWal {
self.failed.len()
}

/// A completed import no longer needs recovery state.
pub(crate) fn remove(&self) -> Result<()> {
if self.dir.exists() {
fs::remove_dir_all(&self.dir)
.with_context(|| format!("remove completed import WAL {}", self.dir.display()))?;
}
Ok(())
}

pub(crate) fn skip_paths(&self) -> HashSet<String> {
self.done
.iter()
Expand Down Expand Up @@ -365,5 +374,8 @@ mod tests {
.unwrap();
assert!(!reset.should_skip("a.json"));
assert_eq!(reset.done_count(), 0);
let dir = reset.dir().to_path_buf();
reset.remove().unwrap();
assert!(!dir.exists());
}
}
53 changes: 22 additions & 31 deletions crates/persisting-pchronicle-cli/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -117,29 +117,19 @@ impl Cli {
}
}

/// Apply S3 backend keys from `--catalog-config` and local `@name` pin settings
/// before the multi-threaded Tokio runtime starts. `std::env::set_var` after
/// worker threads exist is racy on macOS and can leave OpenDAL unable to see
/// `AWS_REGION`.
pub fn apply_catalog_backend_env_before_runtime(cli: &Cli) -> Result<()> {
apply_serve_catalog_backend_env(cli)?;
apply_command_pin_backend_env(cli)?;
Ok(())
/// Run the internal exec worker before constructing a runtime. Ordinary CLI
/// invocations return None and retain their existing execution path.
pub fn run_catalog_worker_before_runtime(cli: &Cli) -> Option<Result<()>> {
match &cli.command {
Command::Serve(args) if args.catalog_query_worker => Some(server::catalog_worker::run()),
_ => None,
}
}

fn apply_serve_catalog_backend_env(cli: &Cli) -> Result<()> {
let Command::Serve(args) = &cli.command else {
return Ok(());
};
if args.command.is_some() || args.catalog_query_worker {
return Ok(());
}
let Some(path) = args.catalog_config.as_ref() else {
return Ok(());
};
let acl = server::catalog::CatalogAcl::load(path)?;
acl.apply_backend_env();
Ok(())
/// Apply local @name pin credentials before creating runtime threads. Catalog
/// server credentials are passed exclusively to isolated workers through IPC.
pub fn apply_catalog_backend_env_before_runtime(cli: &Cli) -> Result<()> {
apply_command_pin_backend_env(cli)
}

fn apply_command_pin_backend_env(cli: &Cli) -> Result<()> {
Expand Down Expand Up @@ -1030,18 +1020,16 @@ struct ServeArgs {
#[arg(long = "gateway-debug", alias = "debug", requires = "gateway_config")]
debug: bool,

/// Directory ACL file (libraries + users). Mounts every [datasets.*] entry
/// into Warehouse and enables catalog:// locators. Mutually exclusive with
/// positional Dataset mounts. Apply S3 endpoint/region/keys from the file
/// before opening stores.
/// Directory ACL file. Authenticate API requests and execute them in
/// bounded, user-scoped worker processes with explicit backend credentials.
#[arg(
long = "catalog-config",
value_name = "FILE",
conflicts_with_all = ["config", "storage", "positional_storage"]
conflicts_with_all = ["config", "storage", "positional_storage", "gateway", "gateway_config", "gateway_dataset", "control"]
)]
catalog_config: Option<PathBuf>,

/// Internal: run one filtered Warehouse request from stdin and exit.
/// Internal: serve framed Warehouse requests over private stdin/stdout IPC.
#[arg(long = "catalog-query-worker", hide = true)]
catalog_query_worker: bool,
}
Expand Down Expand Up @@ -1644,7 +1632,7 @@ pub async fn run_with_stdio(
}) => run_echo(args, &mut diagnostics).await,
Command::Serve(args) => {
if args.catalog_query_worker {
return server::catalog::run_catalog_query_worker().await;
bail!("catalog worker must start before the async runtime");
}
if let Some(command) = args.command {
return run_serve_catalog(command, stdout_is_terminal, stdout, &mut diagnostics);
Expand Down Expand Up @@ -2454,9 +2442,12 @@ fn resolve_serve_config_with_settings(
let storage = serve_storage_uris(args);
let gateway_dataset = resolve_gateway_dataset_uri(args, settings_override)?;
let mut config = if let Some(path) = args.catalog_config.as_deref() {
let acl = server::catalog::CatalogAcl::load(path)?;
acl.apply_backend_env();
server::ChronicleServerConfig::mounted(acl.mounts()?)?
anyhow::ensure!(
gateway_dataset.is_none() && args.control.is_none(),
"catalog workers cannot share Gateway or Control listeners"
);
server::catalog::CatalogAcl::load(path)?;
server::ChronicleServerConfig::front_only()
} else {
match (args.config.as_deref(), storage.as_slice()) {
(Some(config), []) => {
Expand Down
7 changes: 7 additions & 0 deletions crates/persisting-pchronicle-cli/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,13 @@ use persisting_pchronicle_cli::{

fn main() -> ExitCode {
let cli = Cli::parse();
if let Some(result) = persisting_pchronicle_cli::run_catalog_worker_before_runtime(&cli) {
return if result.is_ok() {
ExitCode::SUCCESS
} else {
ExitCode::from(1)
};
}
let debug_errors = cli.debug_errors();
// OpenDAL/Lance read AWS_* from the process environment. Applying catalog
// backend keys after the multi-threaded Tokio runtime starts is racy on
Expand Down
Loading
Loading