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")]
32pub struct PendingEmbeddingsArgs {
33    #[command(subcommand)]
34    pub cmd: PendingEmbeddingsCmd,
35}
36
37#[derive(Debug, Subcommand)]
38pub enum PendingEmbeddingsCmd {
39    /// List all pending embeddings (alias of `embedding list`).
40    List(PendingEmbeddingsListArgs),
41    /// Queue health counts (alias of `embedding status`).
42    Status(crate::commands::embedding::EmbeddingStatusArgs),
43    /// Mark every entry in a given status as abandoned.
44    Abandon(PendingEmbeddingsAbandonArgs),
45}
46
47#[derive(Debug, Args)]
48pub struct PendingEmbeddingsListArgs {
49    /// Filter by status: pending | in_progress | done | abandoned. Default: pending.
50    #[arg(long, default_value = "pending")]
51    pub status: String,
52    /// Maximum number of entries to return. Default: 1000.
53    #[arg(long, default_value_t = 1000)]
54    pub limit: usize,
55    #[arg(long, hide = true)]
56    pub json: bool,
57    #[arg(long)]
58    pub db: Option<String>,
59}
60
61#[derive(Debug, Args)]
62pub struct PendingEmbeddingsAbandonArgs {
63    /// Status to filter: pending | in_progress | done | abandoned. Default: pending.
64    #[arg(long, default_value = "pending")]
65    pub status: String,
66    /// Skip the interactive confirmation prompt.
67    #[arg(long)]
68    pub yes: bool,
69    /// Dry-run: count candidates without modifying.
70    #[arg(long)]
71    pub dry_run: bool,
72    #[arg(long, hide = true)]
73    pub json: bool,
74    #[arg(long)]
75    pub db: Option<String>,
76}
77
78#[derive(Serialize)]
79struct PendingEmbeddingsListEntry {
80    pending_id: i64,
81    memory_id: i64,
82    name: String,
83    namespace: String,
84    backend_chain: String,
85    last_error: Option<String>,
86    last_exit_code: Option<i32>,
87    last_stderr_tail: Option<String>,
88    attempt_count: i32,
89    status: String,
90    updated_at: i64,
91}
92
93impl From<&PendingEmbedding> for PendingEmbeddingsListEntry {
94    fn from(p: &PendingEmbedding) -> Self {
95        Self {
96            pending_id: p.pending_id,
97            memory_id: p.memory_id,
98            name: p.name.clone(),
99            namespace: p.namespace.clone(),
100            backend_chain: p.backend_chain.clone(),
101            last_error: p.last_error.clone(),
102            last_exit_code: p.last_exit_code,
103            last_stderr_tail: p.last_stderr_tail.clone(),
104            attempt_count: p.attempt_count,
105            status: p.status.as_str().to_string(),
106            updated_at: p.updated_at,
107        }
108    }
109}
110
111#[derive(Serialize)]
112struct PendingEmbeddingsListOutput {
113    action: &'static str,
114    filter_status: String,
115    count: usize,
116    entries: Vec<PendingEmbeddingsListEntry>,
117    elapsed_ms: u64,
118}
119
120#[derive(Serialize)]
121struct PendingEmbeddingsAbandonOutput {
122    action: &'static str,
123    dry_run: bool,
124    status: String,
125    candidates: usize,
126    abandoned: usize,
127    elapsed_ms: u64,
128    yes: bool,
129}
130
131pub fn run(args: PendingEmbeddingsArgs) -> Result<(), AppError> {
132    match args.cmd {
133        PendingEmbeddingsCmd::List(a) => run_list(a),
134        PendingEmbeddingsCmd::Status(a) => {
135            crate::commands::embedding::run_status(a, crate::cli::LlmBackendChoice::None)
136        }
137        PendingEmbeddingsCmd::Abandon(a) => run_abandon(a),
138    }
139}
140
141fn parse_status(s: &str) -> Result<PendingEmbeddingStatus, AppError> {
142    match s {
143        "pending" => Ok(PendingEmbeddingStatus::Pending),
144        "in_progress" => Ok(PendingEmbeddingStatus::InProgress),
145        "done" => Ok(PendingEmbeddingStatus::Done),
146        "abandoned" => Ok(PendingEmbeddingStatus::Abandoned),
147        other => Err(AppError::Validation(format!(
148            "invalid status filter: {other} (expected pending|in_progress|done|abandoned)"
149        ))),
150    }
151}
152
153fn open_conn(db: Option<&str>) -> Result<(AppPaths, rusqlite::Connection), AppError> {
154    let paths = AppPaths::resolve(db)?;
155    let conn = open_rw(&paths.db)?;
156    Ok((paths, conn))
157}
158
159fn run_list(args: PendingEmbeddingsListArgs) -> Result<(), AppError> {
160    let start = std::time::Instant::now();
161    let (_paths, conn) = open_conn(args.db.as_deref())?;
162    let status = parse_status(&args.status)?;
163    let rows = pending_embeddings::list_by_status(&conn, status, args.limit)?;
164    let count = rows.len();
165    let entries: Vec<PendingEmbeddingsListEntry> =
166        rows.iter().map(PendingEmbeddingsListEntry::from).collect();
167    let output = PendingEmbeddingsListOutput {
168        action: "pending_embeddings_list",
169        filter_status: status.as_str().to_string(),
170        count,
171        entries,
172        elapsed_ms: start.elapsed().as_millis() as u64,
173    };
174    emit_json_compact(&output)
175}
176
177fn run_abandon(args: PendingEmbeddingsAbandonArgs) -> Result<(), AppError> {
178    let start = std::time::Instant::now();
179    let (_paths, conn) = open_conn(args.db.as_deref())?;
180    let status = parse_status(&args.status)?;
181    let rows = pending_embeddings::list_by_status(&conn, status, 100_000)?;
182    let candidates = rows.len();
183    let mut abandoned = 0usize;
184    if !args.dry_run {
185        for row in &rows {
186            pending_embeddings::abandon(&conn, row.pending_id)?;
187            abandoned += 1;
188        }
189    }
190    let output = PendingEmbeddingsAbandonOutput {
191        action: if args.dry_run {
192            "pending_embeddings_abandon_dry_run"
193        } else {
194            "pending_embeddings_abandon"
195        },
196        dry_run: args.dry_run,
197        status: status.as_str().to_string(),
198        candidates,
199        abandoned,
200        elapsed_ms: start.elapsed().as_millis() as u64,
201        yes: args.yes,
202    };
203    emit_json_compact(&output)
204}