Skip to main content

sqlite_graphrag/storage/
pending_embeddings.rs

1//! GAP-005 (v1.0.82): DAO for the `pending_embeddings` table.
2//!
3//! Queue of memories persisted with a NULL embedding for later reprocessing
4//! via `embedding retry --backend <KIND>` or `enrich --operation re-embed --pending-only`.
5
6use rusqlite::{params, Connection};
7
8use crate::errors::AppError;
9
10/// Pending embedding status.
11#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
12#[serde(rename_all = "snake_case")]
13pub enum PendingEmbeddingStatus {
14    /// Pending variant.
15    Pending,
16    /// In progress variant.
17    InProgress,
18    /// Done variant.
19    Done,
20    /// Abandoned variant.
21    Abandoned,
22}
23
24impl PendingEmbeddingStatus {
25    /// Return the canonical string representation.
26    pub fn as_str(&self) -> &'static str {
27        match self {
28            Self::Pending => "pending",
29            Self::InProgress => "in_progress",
30            Self::Done => "done",
31            Self::Abandoned => "abandoned",
32        }
33    }
34}
35
36/// Pending embedding.
37#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
38pub struct PendingEmbedding {
39    /// Pending ID.
40    pub pending_id: i64,
41    /// Memory identifier.
42    pub memory_id: i64,
43    /// Namespace scope.
44    pub namespace: String,
45    /// Name of this item.
46    pub name: String,
47    /// Backend chain.
48    pub backend_chain: String,
49    /// Last error.
50    pub last_error: Option<String>,
51    /// Last exit code.
52    pub last_exit_code: Option<i32>,
53    /// Last stderr tail.
54    pub last_stderr_tail: Option<String>,
55    /// Attempt count.
56    pub attempt_count: i32,
57    /// Status value.
58    pub status: PendingEmbeddingStatus,
59    /// Creation timestamp.
60    pub created_at: i64,
61    /// Last-update timestamp.
62    pub updated_at: i64,
63}
64
65/// Inserts a new `pending_embeddings` entry with status `pending`.
66#[allow(clippy::too_many_arguments)]
67pub fn insert(
68    conn: &Connection,
69    memory_id: i64,
70    namespace: &str,
71    name: &str,
72    backend_chain: &str,
73    last_error: Option<&str>,
74    last_exit_code: Option<i32>,
75    last_stderr_tail: Option<&str>,
76) -> Result<i64, AppError> {
77    conn.execute(
78        "INSERT INTO pending_embeddings
79            (memory_id, namespace, name, backend_chain, last_error,
80             last_exit_code, last_stderr_tail, attempt_count, status)
81         VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, 0, 'pending')",
82        params![
83            memory_id,
84            namespace,
85            name,
86            backend_chain,
87            last_error,
88            last_exit_code,
89            last_stderr_tail,
90        ],
91    )?;
92    Ok(conn.last_insert_rowid())
93}
94
95/// Update status.
96pub fn update_status(
97    conn: &Connection,
98    pending_id: i64,
99    status: PendingEmbeddingStatus,
100    last_error: Option<&str>,
101    last_exit_code: Option<i32>,
102    last_stderr_tail: Option<&str>,
103) -> Result<(), AppError> {
104    conn.execute(
105        "UPDATE pending_embeddings
106         SET status = ?1,
107             last_error = COALESCE(?2, last_error),
108             last_exit_code = COALESCE(?3, last_exit_code),
109             last_stderr_tail = COALESCE(?4, last_stderr_tail),
110             attempt_count = attempt_count + 1,
111             updated_at = unixepoch()
112         WHERE pending_id = ?5",
113        params![
114            status.as_str(),
115            last_error,
116            last_exit_code,
117            last_stderr_tail,
118            pending_id
119        ],
120    )?;
121    Ok(())
122}
123
124/// List by status.
125pub fn list_by_status(
126    conn: &Connection,
127    status: PendingEmbeddingStatus,
128    limit: usize,
129) -> Result<Vec<PendingEmbedding>, AppError> {
130    let mut stmt = conn.prepare(
131        "SELECT pending_id, memory_id, namespace, name, backend_chain,
132                last_error, last_exit_code, last_stderr_tail,
133                attempt_count, status, created_at, updated_at
134         FROM pending_embeddings
135         WHERE status = ?1
136         ORDER BY updated_at ASC
137         LIMIT ?2",
138    )?;
139    let rows = stmt.query_map(params![status.as_str(), limit as i64], |row| {
140        Ok(PendingEmbedding {
141            pending_id: row.get(0)?,
142            memory_id: row.get(1)?,
143            namespace: row.get(2)?,
144            name: row.get(3)?,
145            backend_chain: row.get(4)?,
146            last_error: row.get(5)?,
147            last_exit_code: row.get(6)?,
148            last_stderr_tail: row.get(7)?,
149            attempt_count: row.get(8)?,
150            status: parse_status(&row.get::<_, String>(9)?).map_err(|e| -> rusqlite::Error {
151                rusqlite::Error::FromSqlConversionFailure(
152                    9,
153                    rusqlite::types::Type::Text,
154                    Box::new(std::io::Error::other(e.to_string())),
155                )
156            })?,
157            created_at: row.get(10)?,
158            updated_at: row.get(11)?,
159        })
160    })?;
161    let mut out = Vec::new();
162    for row in rows {
163        out.push(row?);
164    }
165    Ok(out)
166}
167
168/// Abandon.
169pub fn abandon(conn: &Connection, pending_id: i64) -> Result<(), AppError> {
170    update_status(
171        conn,
172        pending_id,
173        PendingEmbeddingStatus::Abandoned,
174        None,
175        None,
176        None,
177    )
178}
179
180/// Delete.
181pub fn delete(conn: &Connection, pending_id: i64) -> Result<(), AppError> {
182    conn.execute(
183        "DELETE FROM pending_embeddings WHERE pending_id = ?1",
184        params![pending_id],
185    )?;
186    Ok(())
187}
188
189fn parse_status(s: &str) -> Result<PendingEmbeddingStatus, AppError> {
190    match s {
191        "pending" => Ok(PendingEmbeddingStatus::Pending),
192        "in_progress" => Ok(PendingEmbeddingStatus::InProgress),
193        "done" => Ok(PendingEmbeddingStatus::Done),
194        "abandoned" => Ok(PendingEmbeddingStatus::Abandoned),
195        other => Err(AppError::Validation(
196            crate::i18n::validation::unknown_pending_embeddings_status(other),
197        )),
198    }
199}
200
201#[cfg(test)]
202mod tests {
203    use super::*;
204    use rusqlite::Connection;
205
206    fn fresh_db() -> Connection {
207        let mut conn = Connection::open_in_memory().expect("in-memory db");
208        conn.execute_batch("PRAGMA foreign_keys = ON;")
209            .expect("pragma");
210        crate::migrations::runner()
211            .run(&mut conn)
212            .expect("migrations apply");
213        conn
214    }
215
216    fn insert_test_memory(conn: &Connection, name: &str) -> i64 {
217        conn.execute(
218            "INSERT INTO memories (name, namespace, type, description, body, body_hash, source)
219             VALUES (?1, 'global', 'note', 'desc', 'body', 'h', 'agent')",
220            params![name],
221        )
222        .unwrap();
223        conn.last_insert_rowid()
224    }
225
226    #[test]
227    fn insert_records_pending_with_full_diagnostics() {
228        let conn = fresh_db();
229        let mid = insert_test_memory(&conn, "p");
230        let id = insert(
231            &conn,
232            mid,
233            "global",
234            "p",
235            "codex,claude,none",
236            Some("exit 137 SIGKILL"),
237            Some(137),
238            Some("OOM killed by kernel"),
239        )
240        .unwrap();
241        let p = list_by_status(&conn, PendingEmbeddingStatus::Pending, 10)
242            .unwrap()
243            .into_iter()
244            .find(|p| p.pending_id == id)
245            .expect("pending found");
246        assert_eq!(p.backend_chain, "codex,claude,none");
247        assert_eq!(p.last_exit_code, Some(137));
248        assert_eq!(p.last_stderr_tail.as_deref(), Some("OOM killed by kernel"));
249    }
250
251    #[test]
252    fn update_status_increments_attempt_count() {
253        let conn = fresh_db();
254        let mid = insert_test_memory(&conn, "p");
255        let id = insert(&conn, mid, "global", "p", "codex", None, None, None).unwrap();
256        update_status(
257            &conn,
258            id,
259            PendingEmbeddingStatus::InProgress,
260            None,
261            None,
262            None,
263        )
264        .unwrap();
265        let p = list_by_status(&conn, PendingEmbeddingStatus::InProgress, 10)
266            .unwrap()
267            .into_iter()
268            .find(|p| p.pending_id == id)
269            .expect("found");
270        assert_eq!(p.attempt_count, 1);
271    }
272
273    #[test]
274    fn abandon_sets_status() {
275        let conn = fresh_db();
276        let mid = insert_test_memory(&conn, "p");
277        let id = insert(&conn, mid, "global", "p", "codex", None, None, None).unwrap();
278        abandon(&conn, id).unwrap();
279        let abandoned = list_by_status(&conn, PendingEmbeddingStatus::Abandoned, 10).unwrap();
280        assert!(abandoned.iter().any(|p| p.pending_id == id));
281    }
282
283    #[test]
284    fn delete_removes_row() {
285        let conn = fresh_db();
286        let mid = insert_test_memory(&conn, "p");
287        let id = insert(&conn, mid, "global", "p", "codex", None, None, None).unwrap();
288        delete(&conn, id).unwrap();
289        let pending = list_by_status(&conn, PendingEmbeddingStatus::Pending, 10).unwrap();
290        assert!(pending.iter().all(|p| p.pending_id != id));
291    }
292}