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 {
33 #[command(subcommand)]
34 pub cmd: PendingEmbeddingsCmd,
35}
36
37#[derive(Debug, Subcommand)]
38pub enum PendingEmbeddingsCmd {
39 List(PendingEmbeddingsListArgs),
41 Status(crate::commands::embedding::EmbeddingStatusArgs),
43 Abandon(PendingEmbeddingsAbandonArgs),
45}
46
47#[derive(Debug, Args)]
48pub struct PendingEmbeddingsListArgs {
49 #[arg(long, default_value = "pending")]
51 pub status: String,
52 #[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 #[arg(long, default_value = "pending")]
65 pub status: String,
66 #[arg(long)]
68 pub yes: bool,
69 #[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}