Skip to main content

sqlite_graphrag/commands/
pending_embeddings.rs

1//! GAP-005 (v1.0.82): `pending-embeddings` subcommand — high-level batch
2//! operations over the `pending_embeddings` queue.
3//!
4//! ## Subcommands
5//! - `pending-embeddings list` — alias of `embedding list`
6//! - `pending-embeddings retry-all` — bulk re-queue for retry
7//! - `pending-embeddings abandon` — bulk mark abandoned
8//!
9//! The split between `embedding` and `pending-embeddings` mirrors the GAP-005
10//! plan: `embedding` carries per-entry inspection (`status` / `abandon <id>`)
11//! while `pending-embeddings` carries batch operations over the queue as a
12//! whole. The two share the same `pending_embeddings` table and storage
13//! layer.
14
15use clap::{Args, Subcommand};
16use serde::Serialize;
17
18use crate::errors::AppError;
19use crate::output::emit_json_compact;
20use crate::paths::AppPaths;
21use crate::storage::connection::open_rw;
22use crate::storage::pending_embeddings::{self, PendingEmbedding, PendingEmbeddingStatus};
23
24#[derive(Debug, Args)]
25#[command(after_long_help = "EXAMPLES:\n  \
26    # List every pending embedding (alias of `embedding list`)\n  \
27    sqlite-graphrag pending-embeddings list --json\n\n  \
28    # Bulk mark every entry in `pending` status as abandoned\n  \
29    sqlite-graphrag pending-embeddings abandon --status pending --yes\n\n  \
30    # Mark every abandoned entry as abandoned (no-op safe retry)\n  \
31    sqlite-graphrag pending-embeddings abandon --status abandoned --yes")]
32/// Pending embeddings args.
33pub struct PendingEmbeddingsArgs {
34    /// Cmd.
35    #[command(subcommand)]
36    pub cmd: PendingEmbeddingsCmd,
37}
38
39/// Pending embeddings cmd.
40#[derive(Debug, Subcommand)]
41pub enum PendingEmbeddingsCmd {
42    /// List all pending embeddings (alias of `embedding list`).
43    List(PendingEmbeddingsListArgs),
44    /// Queue health counts (alias of `embedding status`).
45    Status(crate::commands::embedding::EmbeddingStatusArgs),
46    /// Mark every entry in a given status as abandoned.
47    Abandon(PendingEmbeddingsAbandonArgs),
48}
49
50/// Pending embeddings list args.
51#[derive(Debug, Args)]
52pub struct PendingEmbeddingsListArgs {
53    /// Filter by status: pending | in_progress | done | abandoned. Default: pending.
54    #[arg(long, default_value = "pending")]
55    pub status: String,
56    /// Maximum number of entries to return. Default: 1000.
57    #[arg(long, default_value_t = 1000)]
58    pub limit: usize,
59    /// Emit machine-readable JSON on stdout.
60    #[arg(long, hide = true)]
61    pub json: bool,
62    /// Path to the SQLite database file.
63    #[arg(long)]
64    pub db: Option<String>,
65}
66
67/// Pending embeddings abandon args.
68#[derive(Debug, Args)]
69pub struct PendingEmbeddingsAbandonArgs {
70    /// Status to filter: pending | in_progress | done | abandoned. Default: pending.
71    #[arg(long, default_value = "pending")]
72    pub status: String,
73    /// Skip the interactive confirmation prompt.
74    #[arg(long)]
75    pub yes: bool,
76    /// Dry-run: count candidates without modifying.
77    #[arg(long)]
78    pub dry_run: bool,
79    /// Emit machine-readable JSON on stdout.
80    #[arg(long, hide = true)]
81    pub json: bool,
82    /// Path to the SQLite database file.
83    #[arg(long)]
84    pub db: Option<String>,
85}
86
87#[derive(Serialize)]
88struct PendingEmbeddingsListEntry {
89    pending_id: i64,
90    memory_id: i64,
91    name: String,
92    namespace: String,
93    backend_chain: String,
94    last_error: Option<String>,
95    last_exit_code: Option<i32>,
96    last_stderr_tail: Option<String>,
97    attempt_count: i32,
98    status: String,
99    updated_at: i64,
100}
101
102impl From<&PendingEmbedding> for PendingEmbeddingsListEntry {
103    fn from(p: &PendingEmbedding) -> Self {
104        Self {
105            pending_id: p.pending_id,
106            memory_id: p.memory_id,
107            name: p.name.clone(),
108            namespace: p.namespace.clone(),
109            backend_chain: p.backend_chain.clone(),
110            last_error: p.last_error.clone(),
111            last_exit_code: p.last_exit_code,
112            last_stderr_tail: p.last_stderr_tail.clone(),
113            attempt_count: p.attempt_count,
114            status: p.status.as_str().to_string(),
115            updated_at: p.updated_at,
116        }
117    }
118}
119
120#[derive(Serialize)]
121struct PendingEmbeddingsListOutput {
122    action: &'static str,
123    filter_status: String,
124    count: usize,
125    entries: Vec<PendingEmbeddingsListEntry>,
126    elapsed_ms: u64,
127}
128
129#[derive(Serialize)]
130struct PendingEmbeddingsAbandonOutput {
131    action: &'static str,
132    dry_run: bool,
133    status: String,
134    candidates: usize,
135    abandoned: usize,
136    elapsed_ms: u64,
137    yes: bool,
138}
139
140/// Run.
141pub fn run(args: PendingEmbeddingsArgs) -> Result<(), AppError> {
142    match args.cmd {
143        PendingEmbeddingsCmd::List(a) => run_list(a),
144        PendingEmbeddingsCmd::Status(a) => {
145            crate::commands::embedding::run_status(a, crate::cli::LlmBackendChoice::None)
146        }
147        PendingEmbeddingsCmd::Abandon(a) => run_abandon(a),
148    }
149}
150
151fn parse_status(s: &str) -> Result<PendingEmbeddingStatus, AppError> {
152    match s {
153        "pending" => Ok(PendingEmbeddingStatus::Pending),
154        "in_progress" => Ok(PendingEmbeddingStatus::InProgress),
155        "done" => Ok(PendingEmbeddingStatus::Done),
156        "abandoned" => Ok(PendingEmbeddingStatus::Abandoned),
157        other => Err(AppError::Validation(crate::i18n::validation::invalid_status_filter(other))),
158    }
159}
160
161fn open_conn(db: Option<&str>) -> Result<(AppPaths, rusqlite::Connection), AppError> {
162    let paths = AppPaths::resolve(db)?;
163    let conn = open_rw(&paths.db)?;
164    Ok((paths, conn))
165}
166
167fn run_list(args: PendingEmbeddingsListArgs) -> Result<(), AppError> {
168    let start = std::time::Instant::now();
169    let (_paths, conn) = open_conn(args.db.as_deref())?;
170    let status = parse_status(&args.status)?;
171    let rows = pending_embeddings::list_by_status(&conn, status, args.limit)?;
172    let count = rows.len();
173    let entries: Vec<PendingEmbeddingsListEntry> =
174        rows.iter().map(PendingEmbeddingsListEntry::from).collect();
175    let output = PendingEmbeddingsListOutput {
176        action: "pending_embeddings_list",
177        filter_status: status.as_str().to_string(),
178        count,
179        entries,
180        elapsed_ms: start.elapsed().as_millis() as u64,
181    };
182    emit_json_compact(&output)
183}
184
185fn run_abandon(args: PendingEmbeddingsAbandonArgs) -> Result<(), AppError> {
186    let start = std::time::Instant::now();
187    let (_paths, conn) = open_conn(args.db.as_deref())?;
188    let status = parse_status(&args.status)?;
189    let rows = pending_embeddings::list_by_status(&conn, status, 100_000)?;
190    let candidates = rows.len();
191    let mut abandoned = 0usize;
192    if !args.dry_run {
193        for row in &rows {
194            pending_embeddings::abandon(&conn, row.pending_id)?;
195            abandoned += 1;
196        }
197    }
198    let output = PendingEmbeddingsAbandonOutput {
199        action: if args.dry_run {
200            "pending_embeddings_abandon_dry_run"
201        } else {
202            "pending_embeddings_abandon"
203        },
204        dry_run: args.dry_run,
205        status: status.as_str().to_string(),
206        candidates,
207        abandoned,
208        elapsed_ms: start.elapsed().as_millis() as u64,
209        yes: args.yes,
210    };
211    emit_json_compact(&output)
212}