Skip to main content

cel_memory_duckdb/
provider.rs

1//! `DuckdbMemoryProvider` — DuckDB embedded backing storage.
2
3use 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
38/// DuckDB-backed [`MemoryProvider`] using `array_cosine_distance` + ILIKE retrieval.
39pub 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    /// Open or create a DuckDB database file, run migrations, and return a provider.
48    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    /// Open an in-memory DuckDB database for tests.
75    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    /// Attach a write hook consulted before every persist.
98    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    /// Attach a summarizer for session summaries and rollups (Phase 3).
104    pub fn with_summarizer(mut self, summarizer: Arc<dyn cel_memory::Summarizer>) -> Self {
105        self.summarizer = Some(summarizer);
106        self
107    }
108
109    /// Connection accessor for integration tests.
110    #[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}