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
45 changes: 45 additions & 0 deletions Cargo.lock

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

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@ blake2 = "0.10.6"
pgvector = { version = "0.4.1", features = ["sqlx", "halfvec"] }
phf = { version = "0.12.1", features = ["macros"] }
indenter = "0.3.4"
indicatif = "0.17.9"
itertools = "0.14.0"
derivative = "2.2.0"
hex = "0.4.3"
Expand Down
61 changes: 59 additions & 2 deletions src/execution/live_updater.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ use crate::{

use super::stats;
use futures::future::try_join_all;
use indicatif::ProgressBar;
use sqlx::PgPool;
use tokio::{sync::watch, task::JoinSet, time::MissedTickBehavior};

Expand Down Expand Up @@ -329,7 +330,47 @@ impl SourceUpdateTask {
update_options: super::source_indexer::UpdateOptions,
) -> Result<()> {
let update_stats = Arc::new(stats::UpdateStats::default());
source_indexing_context

// Spawn periodic stats reporting task if print_stats is enabled
let (reporting_handle, progress_bar) = if self.options.print_stats {
let update_stats_clone = update_stats.clone();
let update_title_owned = update_title.to_string();
let flow_name = self.flow.flow_instance.name.clone();
let import_op_name = self.import_op().name.clone();

// Create a progress bar that will overwrite the same line
let pb = ProgressBar::new_spinner();
pb.set_style(
indicatif::ProgressStyle::default_spinner()
.template("{msg}")
.unwrap(),
);
let pb_clone = pb.clone();

let report_task = async move {
let mut interval = tokio::time::interval(REPORT_INTERVAL);
interval.set_missed_tick_behavior(MissedTickBehavior::Delay);
interval.tick().await; // Skip first tick

loop {
interval.tick().await;
let current_stats = update_stats_clone.as_ref().clone();
if current_stats.has_any_change() {
// Show cumulative stats (always show latest total, not delta)
pb_clone.set_message(format!(
"{}.{} ({update_title_owned}): {}",
flow_name, import_op_name, current_stats
));
}
}
};
(Some(tokio::spawn(report_task)), Some(pb))
} else {
(None, None)
};

// Run the actual update
let update_result = source_indexing_context
.update(&self.pool, &update_stats, update_options)
.await
.with_context(|| {
Expand All @@ -338,12 +379,28 @@ impl SourceUpdateTask {
self.flow.flow_instance.name,
self.import_op().name
)
})?;
});

// Cancel the reporting task if it was spawned
if let Some(handle) = reporting_handle {
handle.abort();
}

// Clear the progress bar to ensure final stats appear on a new line
if let Some(pb) = progress_bar {
pb.finish_and_clear();
}

// Check update result
update_result?;

if update_stats.has_any_change() {
self.status_tx.send_modify(|update| {
update.source_updates_num[self.source_idx] += 1;
});
}

// Report final stats
self.report_stats(&update_stats, update_title);
Ok(())
}
Expand Down
Loading