sqlite_graphrag/commands/enrich/run/
mod.rs1use std::time::Instant;
13
14use rusqlite::Connection;
15
16use super::args::{EnrichArgs, EnrichMode, EnrichOperation};
17use super::events::ConcurrencyEvent;
18use super::queue::{
19 count_eligible_pending, enqueue_candidate, item_type_for, item_type_for_key, open_queue_db,
20 retain_unskipped,
21};
22use super::scheduler;
23use super::DEFAULT_RATE_LIMIT_WAIT;
24use crate::errors::AppError;
25use crate::output::emit_json_line as emit_json;
26use crate::paths::AppPaths;
27use crate::storage::connection::{ensure_db_ready, open_rw};
28
29mod budget;
30mod dry_run;
31mod finalize;
32mod guards;
33mod provider;
34mod queue_prep;
35mod scan_phase;
36
37pub fn run(args: &EnrichArgs, backends: crate::cli::BackendChoice) -> Result<(), AppError> {
39 let crate::cli::BackendChoice {
40 llm: llm_backend,
41 embedding: embedding_backend,
42 } = backends;
43 if guards::handle_pre_db_guards(args, backends)? {
44 return Ok(());
45 }
46
47 let started = Instant::now();
48
49 let paths = AppPaths::resolve(args.db.as_deref())?;
50 ensure_db_ready(&paths)?;
51 let conn = open_rw(&paths.db)?;
52 let namespace = crate::namespace::resolve_namespace(args.namespace.as_deref())?;
53
54 let _singleton = crate::lock::acquire_job_singleton(
60 crate::lock::JobType::Enrich,
61 &namespace,
62 &paths.db,
63 args.wait_job_singleton,
64 args.force_job_singleton,
65 )?;
66
67 let provider_binary = provider::resolve_provider_binary(args)?;
68 provider::check_system_load(args)?;
69 provider::run_preflight(args)?;
70
71 let budget = budget::resolve(args);
72
73 if args.dry_run {
76 let scan = scan_phase::run_scan(&conn, &paths.db, &namespace, args, &budget)?;
77 dry_run::emit_preview(args, &scan.keys, started);
78 return Ok(());
79 }
80
81 let queue_path = crate::paths::sidecar_path(&paths.db, ".enrich-queue.sqlite");
87 let mut queue_conn = open_queue_db(&queue_path)?;
88
89 let op_label = format!("{:?}", args.operation());
92
93 queue_prep::prepare_queue(&queue_conn, &conn, args, &namespace, &op_label, &[])?;
95
96 let item_type = item_type_for(&args.operation());
97 let force_redescribe =
98 args.force_redescribe && matches!(args.operation(), EnrichOperation::EntityDescriptions);
99 let backlog_degree0_proxy = if budget.pair_scan_ops {
101 scan_phase::emit_scan_start(
102 &conn,
103 &namespace,
104 args,
105 &budget,
106 super::events::enrich_operation_cli_name(&args.operation()),
107 )
108 } else {
109 None
110 };
111 let scan_started = Instant::now();
112 let total = super::scan::scan_operation_for_each(&conn, &namespace, args, |page| {
113 if force_redescribe {
114 let _ = queue_prep::reopen_force_redescribe_page(&queue_conn, &namespace, &page);
115 }
116 queue_prep::enqueue_page(
117 &mut queue_conn,
118 &conn,
119 &namespace,
120 &page,
121 item_type,
122 &op_label,
123 )?;
124 Ok(())
125 })?;
126 let scan_elapsed_ms = scan_started.elapsed().as_millis() as u64;
131 emit_json(&super::events::PhaseEvent {
132 phase: "scan",
133 binary_path: None,
134 version: None,
135 items_total: Some(total),
136 items_pending: Some(total),
137 llm_parallelism: args.llm_parallelism,
138 });
139 if budget.pair_scan_ops {
140 let op_cli = super::events::enrich_operation_cli_name(&args.operation());
141 emit_json(&serde_json::json!({
142 "phase": "scan_meta",
143 "operation": op_cli,
144 "pair_algorithm": "cooccurrence+hub_island",
145 "items_total": total,
146 "pairs_enqueued_this_scan": total,
147 "backlog_degree0_proxy": backlog_degree0_proxy,
148 "scan_elapsed_ms": scan_elapsed_ms,
149 "scan_aborted_reason": serde_json::Value::Null,
150 }));
151 }
152 queue_prep::log_enqueue_result(&queue_conn, &op_label, &namespace, total);
153
154 let parallelism = super::events::resolve_drain_parallelism(args);
155
156 let mut counters = super::drain_parallel::DrainCounters::default();
157 let backoff_secs = DEFAULT_RATE_LIMIT_WAIT;
158 let rate_limit_deadline = Instant::now() + crate::runtime_config::rate_limit_deadline_secs();
159 let enrich_started = Instant::now();
160
161 let provider_timeout = match args.mode() {
162 EnrichMode::OpenRouter => args.openrouter_chat_timeout_secs(),
163 };
164
165 let provider_model: Option<&str> = match args.mode() {
166 EnrichMode::OpenRouter => args.openrouter_model.as_deref(),
167 };
168
169 let backoff_clause: &str = if args.ignore_backoff {
173 ""
174 } else {
175 "AND (next_retry_at IS NULL OR next_retry_at <= datetime('now'))"
176 };
177
178 emit_json(&ConcurrencyEvent {
181 phase: "concurrency",
182 scan_parallelism: 1,
183 drain_parallelism: parallelism as u32,
184 });
185
186 let mut until_empty_iter: u32 = 0;
194 let yield_every = scheduler::resolve_yield_every_n(args.yield_every_n_items);
195 let mut yield_count: u64 = 0;
196 let mut items_since_yield: usize = 0;
197 let mut preempted_for_gate = false;
199 loop {
202 if args.until_empty {
203 until_empty_iter = until_empty_iter.saturating_add(1);
204 if until_empty_iter > 1 {
205 let mut rescan = super::events::scan_operation_with_deadline(
209 &conn,
210 &namespace,
211 args,
212 Some(budget.until_deadline),
213 )?;
214 if matches!(args.operation(), EnrichOperation::BodyEnrich) {
222 let _ = retain_unskipped(&queue_conn, &op_label, &mut rescan);
223 }
224 {
226 let tx = queue_conn.transaction()?;
227 let tx_conn: &Connection = &tx;
228 for key in &rescan {
229 let it = item_type_for_key(key, item_type);
230 enqueue_candidate(tx_conn, &conn, &namespace, key, it, &op_label);
231 }
232 tx.commit()?;
233 }
234 }
235 }
236 let completed_before = counters.completed;
237
238 if parallelism > 1 {
242 super::drain_parallel::drain_parallel(
243 super::drain_parallel::ParallelSession {
244 args,
245 paths: &paths,
246 queue_path: &queue_path,
247 namespace: &namespace,
248 parallelism,
249 },
250 super::drain_serial::DrainProvider {
251 binary: provider_binary.as_deref(),
252 model: provider_model,
253 timeout: provider_timeout,
254 backends: crate::cli::BackendChoice::new(llm_backend, embedding_backend),
255 },
256 super::drain_serial::DrainScope {
257 op_label: &op_label,
258 backoff_clause,
259 item_type,
260 total,
261 },
262 &mut counters,
263 )?;
264 } else {
265 super::drain_serial::drain_serial(
266 super::drain_serial::DrainSession {
267 args,
268 paths: &paths,
269 conn: &conn,
270 queue_conn: &queue_conn,
271 namespace: &namespace,
272 },
273 super::drain_serial::DrainProvider {
274 binary: provider_binary.as_deref(),
275 model: provider_model,
276 timeout: provider_timeout,
277 backends: crate::cli::BackendChoice::new(llm_backend, embedding_backend),
278 },
279 super::drain_serial::DrainScope {
280 op_label: &op_label,
281 backoff_clause,
282 item_type,
283 total,
284 },
285 super::drain_serial::DrainClocks {
286 started: enrich_started,
287 until_deadline: budget.until_deadline,
288 rate_limit_deadline,
289 yield_every,
290 backoff_secs,
291 },
292 super::drain_serial::DrainProgress {
293 counters: &mut counters,
294 items_since_yield: &mut items_since_yield,
295 yield_count: &mut yield_count,
296 preempted_for_gate: &mut preempted_for_gate,
297 },
298 )?;
299 }
300
301 if !args.until_empty {
302 break;
303 }
304 let eligible_remaining =
306 count_eligible_pending(&queue_conn, &op_label, &namespace, backoff_clause);
307 let progressed = counters.completed > completed_before;
308 if Instant::now() >= budget.until_deadline {
309 tracing::info!(target: "enrich", "until-empty: max-runtime reached, stopping");
310 break;
311 }
312 if !progressed && eligible_remaining == 0 {
313 tracing::info!(target: "enrich", "until-empty: converged (no eligible items remain)");
314 break;
315 }
316 if eligible_remaining == 0 {
317 std::thread::sleep(std::time::Duration::from_secs(
319 crate::constants::ENRICH_UNTIL_EMPTY_IDLE_NAP_SECS,
320 ));
321 }
322 } finalize::finish(
325 &conn,
326 &queue_conn,
327 &queue_path,
328 args,
329 &op_label,
330 &namespace,
331 finalize::FinalTally {
332 counters: &counters,
333 items_total: total,
334 started,
335 until_deadline: budget.until_deadline,
336 pair_scan_ops: budget.pair_scan_ops,
337 backlog_degree0_proxy,
338 yield_count,
339 preempted_for_gate,
340 },
341 );
342
343 Ok(())
344}