1use std::collections::HashMap;
4use std::path::Path;
5use std::sync::{Arc, Mutex};
6
7use async_trait::async_trait;
8use cel_memory::{
9 CallerScope, MemoryChunk, MemoryError, MemoryProvider, MemoryQuery, MemorySession, MemoryStats,
10 MemoryTier, NewMemoryChunk, NewMemorySession, Result as MemoryResult, SessionFilter,
11 SessionOutcome,
12};
13use chrono::Utc;
14use duckdb::{params, Connection};
15use uuid::Uuid;
16
17use crate::error::DuckdbMemoryError;
18use crate::util::{
19 chunk_matches_query, duck_err, embedding_literal, ids_in_clause, kind_str, optional_row,
20 outcome_str, row_to_chunk, row_to_session, rrf, source_str, tier_str, EMBEDDING_DIM,
21};
22use cel_memory::Embedder;
23
24const GET_CHUNK_SQL: &str =
25 "SELECT id, created_at, kind, tier, source, session_id, project_root, caller_id, content, \
26 CAST(metadata AS VARCHAR) AS metadata, importance, pinned, shareable, superseded_by, \
27 embedding_model, embedding_dim FROM memory_chunks WHERE id = ?";
28
29const SESSION_SELECT: &str =
30 "SELECT id, started_at, ended_at, caller_id, title, summary, outcome, \
31 CAST(metadata AS VARCHAR) AS metadata FROM memory_sessions";
32
33const CHUNK_SELECT: &str =
34 "SELECT id, created_at, kind, tier, source, session_id, project_root, caller_id, content, \
35 CAST(metadata AS VARCHAR) AS metadata, importance, pinned, shareable, superseded_by, \
36 embedding_model, embedding_dim FROM memory_chunks";
37
38pub struct DuckdbMemoryProvider {
40 conn: Arc<Mutex<Connection>>,
41 embedder: Arc<dyn Embedder>,
42 write_hook: Option<Arc<dyn cel_memory::MemoryWriteHook>>,
43 summarizer: Option<Arc<dyn cel_memory::Summarizer>>,
44}
45
46impl DuckdbMemoryProvider {
47 pub async fn open(
49 path: impl AsRef<Path>,
50 embedder: Arc<dyn Embedder>,
51 ) -> Result<Self, DuckdbMemoryError> {
52 let path = path.as_ref().to_path_buf();
53 if embedder.dim() != EMBEDDING_DIM {
54 return Err(DuckdbMemoryError::DimMismatch {
55 expected: EMBEDDING_DIM,
56 actual: embedder.dim(),
57 });
58 }
59 let conn = tokio::task::spawn_blocking(move || -> Result<Connection, DuckdbMemoryError> {
60 let conn = Connection::open(&path)?;
61 crate::migrations::run(&conn)?;
62 Ok(conn)
63 })
64 .await
65 .map_err(|e| DuckdbMemoryError::BlockingJoin(e.to_string()))??;
66 Ok(Self {
67 conn: Arc::new(Mutex::new(conn)),
68 embedder,
69 write_hook: None,
70 summarizer: None,
71 })
72 }
73
74 pub async fn open_in_memory(embedder: Arc<dyn Embedder>) -> Result<Self, DuckdbMemoryError> {
76 if embedder.dim() != EMBEDDING_DIM {
77 return Err(DuckdbMemoryError::DimMismatch {
78 expected: EMBEDDING_DIM,
79 actual: embedder.dim(),
80 });
81 }
82 let conn = tokio::task::spawn_blocking(|| -> Result<Connection, DuckdbMemoryError> {
83 let conn = Connection::open_in_memory()?;
84 crate::migrations::run(&conn)?;
85 Ok(conn)
86 })
87 .await
88 .map_err(|e| DuckdbMemoryError::BlockingJoin(e.to_string()))??;
89 Ok(Self {
90 conn: Arc::new(Mutex::new(conn)),
91 embedder,
92 write_hook: None,
93 summarizer: None,
94 })
95 }
96
97 pub fn with_write_hook(mut self, hook: Arc<dyn cel_memory::MemoryWriteHook>) -> Self {
99 self.write_hook = Some(hook);
100 self
101 }
102
103 pub fn with_summarizer(mut self, summarizer: Arc<dyn cel_memory::Summarizer>) -> Self {
105 self.summarizer = Some(summarizer);
106 self
107 }
108
109 #[doc(hidden)]
111 pub fn conn_for_test(&self) -> Arc<Mutex<Connection>> {
112 Arc::clone(&self.conn)
113 }
114
115 async fn run<F, T>(&self, f: F) -> MemoryResult<T>
116 where
117 F: FnOnce(&Connection) -> MemoryResult<T> + Send + 'static,
118 T: Send + 'static,
119 {
120 let conn = Arc::clone(&self.conn);
121 tokio::task::spawn_blocking(move || {
122 let guard = conn
123 .lock()
124 .map_err(|e| MemoryError::Storage(format!("duckdb lock poisoned: {e}")))?;
125 f(&guard)
126 })
127 .await
128 .map_err(|e| MemoryError::Storage(format!("blocking join failed: {e}")))?
129 }
130}
131
132#[async_trait]
133impl MemoryProvider for DuckdbMemoryProvider {
134 async fn retrieve(&self, query: MemoryQuery) -> MemoryResult<Vec<MemoryChunk>> {
135 if query.text.trim().is_empty() {
136 return Err(MemoryError::InvalidArgument(
137 "query.text must not be empty".into(),
138 ));
139 }
140 let k = query.k.max(1);
141 let candidate_k = (3 * k).max(16);
142 let q_embedding = self.embedder.embed(&query.text).await?;
143 if q_embedding.len() != EMBEDDING_DIM {
144 return Err(MemoryError::Internal(format!(
145 "embedder produced dim {}, expected {EMBEDDING_DIM}",
146 q_embedding.len()
147 )));
148 }
149 let q_lit = embedding_literal(&q_embedding);
150 let scope = query.caller_scope;
151 let caller = query.caller_id.clone();
152 let text = query.text.clone();
153 let query_filter = query.clone();
154
155 self.run(move |conn| {
156 let vector_sql = match scope {
157 CallerScope::Global => format!(
158 "SELECT c.id
159 FROM memory_chunks c
160 INNER JOIN memory_vectors v ON v.chunk_id = c.id
161 ORDER BY array_cosine_distance(v.embedding, {q_lit}::FLOAT[{EMBEDDING_DIM}])
162 LIMIT {candidate_k}"
163 ),
164 CallerScope::Own => format!(
165 "SELECT c.id
166 FROM memory_chunks c
167 INNER JOIN memory_vectors v ON v.chunk_id = c.id
168 WHERE c.caller_id = ?
169 ORDER BY array_cosine_distance(v.embedding, {q_lit}::FLOAT[{EMBEDDING_DIM}])
170 LIMIT {candidate_k}"
171 ),
172 CallerScope::OwnPlusShared => format!(
173 "SELECT c.id
174 FROM memory_chunks c
175 INNER JOIN memory_vectors v ON v.chunk_id = c.id
176 WHERE (c.caller_id = ? OR c.shareable = TRUE)
177 ORDER BY array_cosine_distance(v.embedding, {q_lit}::FLOAT[{EMBEDDING_DIM}])
178 LIMIT {candidate_k}"
179 ),
180 };
181
182 let vector_ids: Vec<String> = match scope {
183 CallerScope::Global => conn
184 .prepare(&vector_sql)
185 .map_err(duck_err)?
186 .query_map([], |row| row.get(0))
187 .map_err(duck_err)?
188 .collect::<Result<_, _>>()
189 .map_err(duck_err)?,
190 CallerScope::Own | CallerScope::OwnPlusShared => conn
191 .prepare(&vector_sql)
192 .map_err(duck_err)?
193 .query_map(params![caller.as_str()], |row| row.get(0))
194 .map_err(duck_err)?
195 .collect::<Result<_, _>>()
196 .map_err(duck_err)?,
197 };
198
199 let lexical_sql = match scope {
200 CallerScope::Global => {
201 "SELECT c.id
202 FROM memory_chunks c
203 WHERE c.content ILIKE '%' || ? || '%'
204 ORDER BY length(c.content)
205 LIMIT ?"
206 }
207 CallerScope::Own => {
208 "SELECT c.id
209 FROM memory_chunks c
210 WHERE c.caller_id = ?
211 AND c.content ILIKE '%' || ? || '%'
212 ORDER BY length(c.content)
213 LIMIT ?"
214 }
215 CallerScope::OwnPlusShared => {
216 "SELECT c.id
217 FROM memory_chunks c
218 WHERE (c.caller_id = ? OR c.shareable = TRUE)
219 AND c.content ILIKE '%' || ? || '%'
220 ORDER BY length(c.content)
221 LIMIT ?"
222 }
223 };
224
225 let lexical_ids: Vec<String> = match scope {
226 CallerScope::Global => conn
227 .prepare(lexical_sql)
228 .map_err(duck_err)?
229 .query_map(params![text.as_str(), candidate_k as i64], |row| row.get(0))
230 .map_err(duck_err)?
231 .collect::<Result<_, _>>()
232 .map_err(duck_err)?,
233 CallerScope::Own | CallerScope::OwnPlusShared => conn
234 .prepare(lexical_sql)
235 .map_err(duck_err)?
236 .query_map(
237 params![caller.as_str(), text.as_str(), candidate_k as i64],
238 |row| row.get(0),
239 )
240 .map_err(duck_err)?
241 .collect::<Result<_, _>>()
242 .map_err(duck_err)?,
243 };
244
245 let mut scores: HashMap<String, f32> = HashMap::new();
246 for (rank, id) in vector_ids.into_iter().enumerate() {
247 *scores.entry(id).or_insert(0.0) += rrf(rank, 60.0);
248 }
249 for (rank, id) in lexical_ids.into_iter().enumerate() {
250 *scores.entry(id).or_insert(0.0) += rrf(rank, 60.0);
251 }
252 if scores.is_empty() {
253 return Ok(Vec::new());
254 }
255
256 let ids: Vec<String> = scores.keys().cloned().collect();
257 let select_sql = format!(
258 "{CHUNK_SELECT} WHERE id IN ({})",
259 ids_in_clause(&ids)
260 );
261 let mut stmt = conn.prepare(&select_sql).map_err(duck_err)?;
262 let rows = stmt
263 .query_map([], row_to_chunk)
264 .map_err(duck_err)?;
265
266 let mut filtered = Vec::new();
267 for row in rows {
268 let c = row.map_err(duck_err)?;
269 if chunk_matches_query(&c, &query_filter) {
270 filtered.push(c);
271 }
272 }
273 filtered.sort_by(|a, b| {
274 let sa = scores.get(&a.id).copied().unwrap_or(0.0);
275 let sb = scores.get(&b.id).copied().unwrap_or(0.0);
276 sb.partial_cmp(&sa).unwrap_or(std::cmp::Ordering::Equal)
277 });
278 filtered.truncate(k);
279 Ok(filtered)
280 })
281 .await
282 }
283
284 async fn get(&self, chunk_id: &str) -> MemoryResult<Option<MemoryChunk>> {
285 let chunk_id = chunk_id.to_string();
286 self.run(move |conn| {
287 optional_row(
288 conn.query_row(GET_CHUNK_SQL, params![chunk_id.as_str()], row_to_chunk),
289 Ok,
290 )
291 })
292 .await
293 }
294
295 async fn get_session(&self, session_id: &str) -> MemoryResult<Option<MemorySession>> {
296 let session_id = session_id.to_string();
297 self.run(move |conn| {
298 optional_row(
299 conn.query_row(
300 &format!("{SESSION_SELECT} WHERE id = ?"),
301 params![session_id.as_str()],
302 row_to_session,
303 ),
304 Ok,
305 )
306 })
307 .await
308 }
309
310 async fn list_sessions(&self, filter: SessionFilter) -> MemoryResult<Vec<MemorySession>> {
311 self.run(move |conn| list_sessions_inner(conn, filter))
312 .await
313 }
314
315 async fn write(&self, new_chunk: NewMemoryChunk) -> MemoryResult<MemoryChunk> {
316 if new_chunk.content.trim().is_empty() {
317 return Err(MemoryError::InvalidArgument(
318 "content must not be empty".into(),
319 ));
320 }
321 if let Some(hook) = &self.write_hook {
322 match hook.before_write(&new_chunk).await? {
323 cel_memory::WriteDecision::Allow => {}
324 cel_memory::WriteDecision::Redact { reason } => {
325 return Ok(MemoryChunk {
326 id: Uuid::now_v7().to_string(),
327 created_at: Utc::now(),
328 kind: new_chunk.kind,
329 tier: MemoryTier::Session,
330 source: new_chunk.source,
331 session_id: new_chunk.session_id,
332 project_root: new_chunk.project_root,
333 caller_id: new_chunk.caller_id,
334 content: format!("<redacted: {reason}>"),
335 metadata: serde_json::json!({"redacted": true, "reason": reason}),
336 importance: 0.0,
337 pinned: false,
338 shareable: false,
339 superseded_by: None,
340 embedding_model: "none".into(),
341 embedding_dim: 0,
342 });
343 }
344 }
345 }
346
347 let embedding = self.embedder.embed(&new_chunk.content).await?;
348 if embedding.len() != EMBEDDING_DIM {
349 return Err(MemoryError::Internal(format!(
350 "embedder produced dim {}, expected {EMBEDDING_DIM}",
351 embedding.len()
352 )));
353 }
354
355 let chunk = MemoryChunk {
356 id: Uuid::now_v7().to_string(),
357 created_at: Utc::now(),
358 kind: new_chunk.kind,
359 tier: MemoryTier::Session,
360 source: new_chunk.source,
361 session_id: new_chunk.session_id.clone(),
362 project_root: new_chunk.project_root.clone(),
363 caller_id: new_chunk.caller_id.clone(),
364 content: new_chunk.content.clone(),
365 metadata: new_chunk.metadata.clone(),
366 importance: cel_memory::score_importance(&new_chunk),
367 pinned: new_chunk.pinned,
368 shareable: new_chunk.shareable,
369 superseded_by: None,
370 embedding_model: self.embedder.model_name().to_string(),
371 embedding_dim: EMBEDDING_DIM as u32,
372 };
373
374 let metadata_json = if chunk.metadata.is_null() {
375 "{}".to_string()
376 } else {
377 chunk.metadata.to_string()
378 };
379 let emb_lit = embedding_literal(&embedding);
380 let chunk_id = chunk.id.clone();
381
382 self.run(move |conn| {
383 conn.execute_batch("BEGIN TRANSACTION")
384 .map_err(|e| MemoryError::Storage(e.to_string()))?;
385 let result = (|| -> MemoryResult<()> {
386 conn.execute(
387 "INSERT INTO memory_chunks(
388 id, created_at, kind, tier, source, session_id, project_root,
389 caller_id, content, metadata, importance, pinned, shareable,
390 superseded_by, embedding_model, embedding_dim
391 ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
392 params![
393 chunk.id.as_str(),
394 chunk.created_at,
395 kind_str(chunk.kind),
396 tier_str(chunk.tier),
397 source_str(chunk.source),
398 chunk.session_id.as_deref(),
399 chunk.project_root.as_deref(),
400 chunk.caller_id.as_str(),
401 chunk.content.as_str(),
402 metadata_json.as_str(),
403 chunk.importance,
404 chunk.pinned,
405 chunk.shareable,
406 chunk.superseded_by.as_deref(),
407 chunk.embedding_model.as_str(),
408 chunk.embedding_dim as i32,
409 ],
410 )
411 .map_err(|e| MemoryError::Storage(e.to_string()))?;
412 conn.execute_batch(&format!(
413 "INSERT INTO memory_vectors (chunk_id, embedding) VALUES ({}, {}::FLOAT[{EMBEDDING_DIM}])",
414 crate::util::sql_string_literal(&chunk_id),
415 emb_lit
416 ))
417 .map_err(|e| MemoryError::Storage(e.to_string()))?;
418 Ok(())
419 })();
420 if result.is_ok() {
421 conn.execute_batch("COMMIT")
422 .map_err(|e| MemoryError::Storage(e.to_string()))?;
423 } else {
424 conn.execute_batch("ROLLBACK").ok();
425 result?;
426 }
427 Ok(chunk)
428 })
429 .await
430 }
431
432 async fn write_batch(&self, chunks: Vec<NewMemoryChunk>) -> MemoryResult<Vec<MemoryChunk>> {
433 let mut out = Vec::with_capacity(chunks.len());
434 for chunk in chunks {
435 out.push(self.write(chunk).await?);
436 }
437 Ok(out)
438 }
439
440 async fn open_session(&self, init: NewMemorySession) -> MemoryResult<MemorySession> {
441 let session = MemorySession {
442 id: Uuid::now_v7().to_string(),
443 started_at: Utc::now(),
444 ended_at: None,
445 caller_id: init.caller_id,
446 title: init.title,
447 summary: None,
448 outcome: SessionOutcome::Open,
449 metadata: init.metadata,
450 };
451 let metadata_json = if session.metadata.is_null() {
452 "{}".to_string()
453 } else {
454 session.metadata.to_string()
455 };
456 let session_id = session.id.clone();
457 self.run(move |conn| {
458 conn.execute(
459 "INSERT INTO memory_sessions (id, started_at, ended_at, caller_id, title, summary, outcome, metadata)
460 VALUES (?, ?, ?, ?, ?, ?, ?, ?)",
461 params![
462 session_id.as_str(),
463 session.started_at,
464 session.ended_at,
465 session.caller_id.as_str(),
466 session.title.as_deref(),
467 session.summary.as_deref(),
468 outcome_str(session.outcome),
469 metadata_json.as_str(),
470 ],
471 )
472 .map_err(|e| MemoryError::Storage(e.to_string()))?;
473 Ok(session)
474 })
475 .await
476 }
477
478 async fn close_session(&self, session_id: &str, outcome: SessionOutcome) -> MemoryResult<()> {
479 let session_id = session_id.to_string();
480 self.run(move |conn| {
481 let changed = conn
482 .execute(
483 "UPDATE memory_sessions SET ended_at = ?, outcome = ? WHERE id = ?",
484 params![Utc::now(), outcome_str(outcome), session_id.as_str()],
485 )
486 .map_err(|e| MemoryError::Storage(e.to_string()))?;
487 if changed == 0 {
488 return Err(MemoryError::NotFound(session_id));
489 }
490 Ok(())
491 })
492 .await
493 }
494
495 async fn rename_session(&self, session_id: &str, title: &str) -> MemoryResult<()> {
496 let session_id = session_id.to_string();
497 let title = title.to_string();
498 self.run(move |conn| {
499 let changed = conn
500 .execute(
501 "UPDATE memory_sessions SET title = ? WHERE id = ?",
502 params![title.as_str(), session_id.as_str()],
503 )
504 .map_err(|e| MemoryError::Storage(e.to_string()))?;
505 if changed == 0 {
506 return Err(MemoryError::NotFound(session_id));
507 }
508 Ok(())
509 })
510 .await
511 }
512
513 async fn stats(&self) -> MemoryResult<MemoryStats> {
514 self.run(|conn| {
515 let total_chunks: i64 = conn
516 .query_row("SELECT COUNT(*)::BIGINT FROM memory_chunks", [], |row| {
517 row.get(0)
518 })
519 .map_err(|e| MemoryError::Storage(e.to_string()))?;
520 let total_sessions: i64 = conn
521 .query_row("SELECT COUNT(*)::BIGINT FROM memory_sessions", [], |row| {
522 row.get(0)
523 })
524 .map_err(|e| MemoryError::Storage(e.to_string()))?;
525 let embedding_model = optional_row(
526 conn.query_row(
527 "SELECT embedding_model FROM memory_chunks ORDER BY created_at DESC LIMIT 1",
528 [],
529 |row| row.get::<_, String>(0),
530 ),
531 Ok,
532 )?;
533 Ok(MemoryStats {
534 total_chunks: total_chunks as usize,
535 total_sessions: total_sessions as usize,
536 embedding_model,
537 ..MemoryStats::default()
538 })
539 })
540 .await
541 }
542
543 async fn summarize_session(&self, _session_id: &str) -> MemoryResult<MemoryChunk> {
544 if self.summarizer.is_none() {
545 return Err(MemoryError::NotImplemented(
546 "DuckdbMemoryProvider::summarize_session — attach summarizer via with_summarizer",
547 ));
548 }
549 Err(MemoryError::NotImplemented(
550 "DuckdbMemoryProvider::summarize_session — Phase 3",
551 ))
552 }
553
554 async fn rollup_day(&self, _date: chrono::NaiveDate) -> MemoryResult<Vec<MemoryChunk>> {
555 Err(MemoryError::NotImplemented(
556 "DuckdbMemoryProvider::rollup_day — Phase 3",
557 ))
558 }
559
560 async fn rollup_rule_week(
561 &self,
562 _rule_id: &str,
563 _week_start: chrono::NaiveDate,
564 ) -> MemoryResult<MemoryChunk> {
565 Err(MemoryError::NotImplemented(
566 "DuckdbMemoryProvider::rollup_rule_week — Phase 3",
567 ))
568 }
569
570 async fn run_aging_sweep(&self) -> MemoryResult<cel_memory::AgingReport> {
571 Err(MemoryError::NotImplemented(
572 "DuckdbMemoryProvider::run_aging_sweep — Phase 2",
573 ))
574 }
575
576 async fn re_embed_all(&self, _target_model: &str) -> MemoryResult<cel_memory::ReEmbedReport> {
577 Err(MemoryError::NotImplemented(
578 "DuckdbMemoryProvider::re_embed_all — Phase 4",
579 ))
580 }
581
582 async fn export(
583 &self,
584 _filter: cel_memory::ExportFilter,
585 ) -> MemoryResult<cel_memory::ExportBundle> {
586 Err(MemoryError::NotImplemented(
587 "DuckdbMemoryProvider::export — Phase 2",
588 ))
589 }
590
591 async fn pin(&self, _chunk_id: &str, _pinned: bool) -> MemoryResult<()> {
592 Err(MemoryError::NotImplemented(
593 "DuckdbMemoryProvider::pin — Phase 2",
594 ))
595 }
596
597 async fn update_importance(&self, _chunk_id: &str, _importance: f32) -> MemoryResult<()> {
598 Err(MemoryError::NotImplemented(
599 "DuckdbMemoryProvider::update_importance — Phase 2",
600 ))
601 }
602
603 async fn supersede(&self, _old_id: &str, _new_id: &str) -> MemoryResult<()> {
604 Err(MemoryError::NotImplemented(
605 "DuckdbMemoryProvider::supersede — Phase 2",
606 ))
607 }
608
609 async fn record_access(
610 &self,
611 _chunk_id: &str,
612 _retrieved_by: &str,
613 _used: bool,
614 ) -> MemoryResult<()> {
615 Err(MemoryError::NotImplemented(
616 "DuckdbMemoryProvider::record_access — Phase 2",
617 ))
618 }
619
620 async fn delete(
621 &self,
622 _chunk_id: &str,
623 _reason: cel_memory::EvictionReason,
624 ) -> MemoryResult<()> {
625 Err(MemoryError::NotImplemented(
626 "DuckdbMemoryProvider::delete — Phase 2",
627 ))
628 }
629
630 async fn delete_matching(
631 &self,
632 _predicate: cel_memory::MemoryPredicate,
633 _reason: cel_memory::EvictionReason,
634 ) -> MemoryResult<usize> {
635 Err(MemoryError::NotImplemented(
636 "DuckdbMemoryProvider::delete_matching — Phase 2",
637 ))
638 }
639
640 async fn purge_all(&self) -> MemoryResult<cel_memory::PurgeReport> {
641 Err(MemoryError::NotImplemented(
642 "DuckdbMemoryProvider::purge_all — Phase 2",
643 ))
644 }
645}
646
647fn list_sessions_inner(
648 conn: &Connection,
649 filter: SessionFilter,
650) -> MemoryResult<Vec<MemorySession>> {
651 let mut out = Vec::new();
652 match (&filter.caller_id, filter.open_only, filter.outcome) {
653 (None, false, None) => {
654 let mut stmt = conn
655 .prepare(&format!("{SESSION_SELECT} ORDER BY started_at DESC"))
656 .map_err(duck_err)?;
657 for row in stmt
658 .query_map([], row_to_session)
659 .map_err(duck_err)? {
660 out.push(row.map_err(duck_err)?);
661 }
662 }
663 (Some(caller), false, None) => {
664 let mut stmt = conn
665 .prepare(&format!(
666 "{SESSION_SELECT} WHERE caller_id = ? ORDER BY started_at DESC"
667 ))
668 .map_err(duck_err)?;
669 for row in stmt
670 .query_map(params![caller.as_str()], row_to_session)
671 .map_err(duck_err)?
672 {
673 out.push(row.map_err(duck_err)?);
674 }
675 }
676 (None, true, _) => {
677 let mut stmt = conn
678 .prepare(&format!(
679 "{SESSION_SELECT} WHERE outcome = 'open' ORDER BY started_at DESC"
680 ))
681 .map_err(duck_err)?;
682 for row in stmt
683 .query_map([], row_to_session)
684 .map_err(duck_err)? {
685 out.push(row.map_err(duck_err)?);
686 }
687 }
688 (Some(caller), true, _) => {
689 let mut stmt = conn
690 .prepare(&format!(
691 "{SESSION_SELECT} WHERE caller_id = ? AND outcome = 'open' ORDER BY started_at DESC"
692 ))
693 .map_err(duck_err)?;
694 for row in stmt
695 .query_map(params![caller.as_str()], row_to_session)
696 .map_err(duck_err)?
697 {
698 out.push(row.map_err(duck_err)?);
699 }
700 }
701 (None, false, Some(outcome)) => {
702 let mut stmt = conn
703 .prepare(&format!(
704 "{SESSION_SELECT} WHERE outcome = ? ORDER BY started_at DESC"
705 ))
706 .map_err(duck_err)?;
707 for row in stmt
708 .query_map(params![outcome_str(outcome)], row_to_session)
709 .map_err(duck_err)?
710 {
711 out.push(row.map_err(duck_err)?);
712 }
713 }
714 (Some(caller), false, Some(outcome)) => {
715 let mut stmt = conn
716 .prepare(&format!(
717 "{SESSION_SELECT} WHERE caller_id = ? AND outcome = ? ORDER BY started_at DESC"
718 ))
719 .map_err(duck_err)?;
720 for row in stmt
721 .query_map(params![caller.as_str(), outcome_str(outcome)], row_to_session)
722 .map_err(duck_err)?
723 {
724 out.push(row.map_err(duck_err)?);
725 }
726 }
727 }
728 Ok(out)
729}