1use rusqlite::{params, Connection};
7
8use crate::errors::AppError;
9
10#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
12#[serde(rename_all = "snake_case")]
13pub enum PendingEmbeddingStatus {
14 Pending,
16 InProgress,
18 Done,
20 Abandoned,
22}
23
24impl PendingEmbeddingStatus {
25 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#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
38pub struct PendingEmbedding {
39 pub pending_id: i64,
41 pub memory_id: i64,
43 pub namespace: String,
45 pub name: String,
47 pub backend_chain: String,
49 pub last_error: Option<String>,
51 pub last_exit_code: Option<i32>,
53 pub last_stderr_tail: Option<String>,
55 pub attempt_count: i32,
57 pub status: PendingEmbeddingStatus,
59 pub created_at: i64,
61 pub updated_at: i64,
63}
64
65#[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
95pub 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
124pub 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
168pub 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
180pub 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(crate::i18n::validation::unknown_pending_embeddings_status(other))),
196 }
197}
198
199#[cfg(test)]
200mod tests {
201 use super::*;
202 use rusqlite::Connection;
203
204 fn fresh_db() -> Connection {
205 let mut conn = Connection::open_in_memory().expect("in-memory db");
206 conn.execute_batch("PRAGMA foreign_keys = ON;")
207 .expect("pragma");
208 crate::migrations::runner()
209 .run(&mut conn)
210 .expect("migrations apply");
211 conn
212 }
213
214 fn insert_test_memory(conn: &Connection, name: &str) -> i64 {
215 conn.execute(
216 "INSERT INTO memories (name, namespace, type, description, body, body_hash, source)
217 VALUES (?1, 'global', 'note', 'desc', 'body', 'h', 'agent')",
218 params![name],
219 )
220 .unwrap();
221 conn.last_insert_rowid()
222 }
223
224 #[test]
225 fn insert_records_pending_with_full_diagnostics() {
226 let conn = fresh_db();
227 let mid = insert_test_memory(&conn, "p");
228 let id = insert(
229 &conn,
230 mid,
231 "global",
232 "p",
233 "codex,claude,none",
234 Some("exit 137 SIGKILL"),
235 Some(137),
236 Some("OOM killed by kernel"),
237 )
238 .unwrap();
239 let p = list_by_status(&conn, PendingEmbeddingStatus::Pending, 10)
240 .unwrap()
241 .into_iter()
242 .find(|p| p.pending_id == id)
243 .expect("pending found");
244 assert_eq!(p.backend_chain, "codex,claude,none");
245 assert_eq!(p.last_exit_code, Some(137));
246 assert_eq!(p.last_stderr_tail.as_deref(), Some("OOM killed by kernel"));
247 }
248
249 #[test]
250 fn update_status_increments_attempt_count() {
251 let conn = fresh_db();
252 let mid = insert_test_memory(&conn, "p");
253 let id = insert(&conn, mid, "global", "p", "codex", None, None, None).unwrap();
254 update_status(
255 &conn,
256 id,
257 PendingEmbeddingStatus::InProgress,
258 None,
259 None,
260 None,
261 )
262 .unwrap();
263 let p = list_by_status(&conn, PendingEmbeddingStatus::InProgress, 10)
264 .unwrap()
265 .into_iter()
266 .find(|p| p.pending_id == id)
267 .expect("found");
268 assert_eq!(p.attempt_count, 1);
269 }
270
271 #[test]
272 fn abandon_sets_status() {
273 let conn = fresh_db();
274 let mid = insert_test_memory(&conn, "p");
275 let id = insert(&conn, mid, "global", "p", "codex", None, None, None).unwrap();
276 abandon(&conn, id).unwrap();
277 let abandoned = list_by_status(&conn, PendingEmbeddingStatus::Abandoned, 10).unwrap();
278 assert!(abandoned.iter().any(|p| p.pending_id == id));
279 }
280
281 #[test]
282 fn delete_removes_row() {
283 let conn = fresh_db();
284 let mid = insert_test_memory(&conn, "p");
285 let id = insert(&conn, mid, "global", "p", "codex", None, None, None).unwrap();
286 delete(&conn, id).unwrap();
287 let pending = list_by_status(&conn, PendingEmbeddingStatus::Pending, 10).unwrap();
288 assert!(pending.iter().all(|p| p.pending_id != id));
289 }
290}