sqlite_graphrag/commands/
pending_embeddings.rs1use 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 {
34 #[command(subcommand)]
36 pub cmd: PendingEmbeddingsCmd,
37}
38
39#[derive(Debug, Subcommand)]
41pub enum PendingEmbeddingsCmd {
42 List(PendingEmbeddingsListArgs),
44 Status(crate::commands::embedding::EmbeddingStatusArgs),
46 Abandon(PendingEmbeddingsAbandonArgs),
48}
49
50#[derive(Debug, Args)]
52pub struct PendingEmbeddingsListArgs {
53 #[arg(long, default_value = "pending")]
55 pub status: String,
56 #[arg(long, default_value_t = 1000)]
58 pub limit: usize,
59 #[arg(long, hide = true)]
61 pub json: bool,
62 #[arg(long)]
64 pub db: Option<String>,
65}
66
67#[derive(Debug, Args)]
69pub struct PendingEmbeddingsAbandonArgs {
70 #[arg(long, default_value = "pending")]
72 pub status: String,
73 #[arg(long)]
75 pub yes: bool,
76 #[arg(long)]
78 pub dry_run: bool,
79 #[arg(long, hide = true)]
81 pub json: bool,
82 #[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
140pub 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(
158 crate::i18n::validation::invalid_status_filter(other),
159 )),
160 }
161}
162
163fn open_conn(db: Option<&str>) -> Result<(AppPaths, rusqlite::Connection), AppError> {
164 let paths = AppPaths::resolve(db)?;
165 let conn = open_rw(&paths.db)?;
166 Ok((paths, conn))
167}
168
169fn run_list(args: PendingEmbeddingsListArgs) -> Result<(), AppError> {
170 let start = std::time::Instant::now();
171 let (_paths, conn) = open_conn(args.db.as_deref())?;
172 let status = parse_status(&args.status)?;
173 let rows = pending_embeddings::list_by_status(&conn, status, args.limit)?;
174 let count = rows.len();
175 let entries: Vec<PendingEmbeddingsListEntry> =
176 rows.iter().map(PendingEmbeddingsListEntry::from).collect();
177 let output = PendingEmbeddingsListOutput {
178 action: "pending_embeddings_list",
179 filter_status: status.as_str().to_string(),
180 count,
181 entries,
182 elapsed_ms: start.elapsed().as_millis() as u64,
183 };
184 emit_json_compact(&output)
185}
186
187fn run_abandon(args: PendingEmbeddingsAbandonArgs) -> Result<(), AppError> {
188 let start = std::time::Instant::now();
189 let (_paths, conn) = open_conn(args.db.as_deref())?;
190 let status = parse_status(&args.status)?;
191 let rows = pending_embeddings::list_by_status(&conn, status, 100_000)?;
192 let candidates = rows.len();
193 let mut abandoned = 0usize;
194 if !args.dry_run {
195 for row in &rows {
196 pending_embeddings::abandon(&conn, row.pending_id)?;
197 abandoned += 1;
198 }
199 }
200 let output = PendingEmbeddingsAbandonOutput {
201 action: if args.dry_run {
202 "pending_embeddings_abandon_dry_run"
203 } else {
204 "pending_embeddings_abandon"
205 },
206 dry_run: args.dry_run,
207 status: status.as_str().to_string(),
208 candidates,
209 abandoned,
210 elapsed_ms: start.elapsed().as_millis() as u64,
211 yes: args.yes,
212 };
213 emit_json_compact(&output)
214}