Skip to main content

aion_store_libsql/
visibility.rs

1//! Workflow visibility projection storage backed by libSQL.
2
3use std::collections::HashMap;
4use std::fmt::Write as _;
5
6use aion_core::{RunId, SearchAttributeValue, WorkflowId, WorkflowStatus};
7use aion_store::StoreError;
8use aion_store::visibility::{
9    ListWorkflowsFilter, SearchAttributePredicate, VisibilityRecord, VisibilityStore,
10    WorkflowSummary as VisibilityWorkflowSummary,
11};
12use async_trait::async_trait;
13use chrono::{DateTime, SecondsFormat, Utc};
14use libsql::{Value, params_from_iter};
15use uuid::Uuid;
16
17use crate::store::LibSqlStore;
18
19const UPSERT_VISIBILITY_SQL: &str = "
20INSERT OR REPLACE INTO visibility (
21    workflow_id,
22    run_id,
23    workflow_type,
24    status,
25    start_time,
26    close_time,
27    failed_step,
28    failure_reason,
29    search_attributes
30)
31VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)";
32
33/// Upsert a complete workflow visibility projection row.
34///
35/// # Errors
36///
37/// Returns `StoreError::Serialization` when record fields cannot be encoded and
38/// `StoreError::Backend` when libSQL rejects the upsert.
39pub(crate) async fn record_visibility(
40    conn: &libsql::Connection,
41    record: VisibilityRecord,
42) -> Result<(), StoreError> {
43    let workflow_id = record.workflow_id.to_string();
44    let run_id = record.run_id.to_string();
45    let status = encode_status(record.status)?;
46    let start_time = encode_timestamp(record.start_time);
47    let close_time = record.close_time.map(encode_timestamp);
48    let failed_step = record.failed_step;
49    let failure_reason = record.failure_reason;
50    let search_attributes = serde_json::to_string(&record.search_attributes)
51        .map_err(|error| crate::error::serde_json_error(&error))?;
52
53    conn.execute(
54        UPSERT_VISIBILITY_SQL,
55        (
56            workflow_id,
57            run_id,
58            record.workflow_type,
59            status,
60            start_time,
61            close_time,
62            failed_step,
63            failure_reason,
64            search_attributes,
65        ),
66    )
67    .await
68    .map_err(|error| crate::error::libsql_error(&error))?;
69
70    Ok(())
71}
72
73/// List workflow visibility summaries matching `filter`.
74///
75/// # Errors
76///
77/// Returns backend errors from libSQL or serialization errors when persisted rows cannot be decoded.
78pub(crate) async fn list_workflows(
79    conn: &libsql::Connection,
80    filter: ListWorkflowsFilter,
81) -> Result<Vec<VisibilityWorkflowSummary>, StoreError> {
82    let plan = QueryPlan::list(&filter)?;
83    let mut rows = conn
84        .query(&plan.sql, params_from_iter(plan.params))
85        .await
86        .map_err(|error| crate::error::libsql_error(&error))?;
87
88    let mut summaries = Vec::new();
89    while let Some(row) = rows
90        .next()
91        .await
92        .map_err(|error| crate::error::libsql_error(&error))?
93    {
94        summaries.push(decode_summary(&row)?);
95    }
96
97    Ok(summaries)
98}
99
100/// Count workflow visibility summaries matching `filter`, ignoring pagination fields.
101///
102/// # Errors
103///
104/// Returns backend errors from libSQL or serialization errors from filter encoding.
105pub(crate) async fn count_workflows(
106    conn: &libsql::Connection,
107    filter: ListWorkflowsFilter,
108) -> Result<u64, StoreError> {
109    let plan = QueryPlan::count(&filter)?;
110    let mut rows = conn
111        .query(&plan.sql, params_from_iter(plan.params))
112        .await
113        .map_err(|error| crate::error::libsql_error(&error))?;
114    let row = rows
115        .next()
116        .await
117        .map_err(|error| crate::error::libsql_error(&error))?
118        .ok_or_else(|| {
119            StoreError::Backend(String::from("visibility count query returned no row"))
120        })?;
121    let count: i64 = row
122        .get(0)
123        .map_err(|error| crate::error::libsql_error(&error))?;
124
125    u64::try_from(count).map_err(|error| StoreError::Backend(error.to_string()))
126}
127
128#[async_trait]
129impl VisibilityStore for LibSqlStore {
130    async fn record_visibility(&self, record: VisibilityRecord) -> Result<(), StoreError> {
131        record_visibility(self.connection(), record).await
132    }
133
134    async fn list_workflows(
135        &self,
136        filter: ListWorkflowsFilter,
137    ) -> Result<Vec<VisibilityWorkflowSummary>, StoreError> {
138        list_workflows(self.connection(), filter).await
139    }
140
141    async fn count_workflows(&self, filter: ListWorkflowsFilter) -> Result<u64, StoreError> {
142        count_workflows(self.connection(), filter).await
143    }
144}
145
146struct QueryPlan {
147    sql: String,
148    params: Vec<Value>,
149}
150
151impl QueryPlan {
152    fn list(filter: &ListWorkflowsFilter) -> Result<Self, StoreError> {
153        let mut plan = Self::filtered(
154            "SELECT workflow_id, run_id, workflow_type, status, start_time, close_time, failed_step, failure_reason, search_attributes FROM visibility",
155            filter,
156        )?;
157        plan.sql
158            .push_str(" ORDER BY start_time DESC, workflow_id ASC");
159        if let Some(limit) = filter.limit {
160            plan.push_param(Value::Integer(i64::from(limit)));
161            write!(plan.sql, " LIMIT ?{}", plan.params.len())
162                .map_err(|error| StoreError::Backend(error.to_string()))?;
163        }
164        if let Some(offset) = filter.offset {
165            plan.push_param(Value::Integer(i64::from(offset)));
166            write!(plan.sql, " OFFSET ?{}", plan.params.len())
167                .map_err(|error| StoreError::Backend(error.to_string()))?;
168        }
169        Ok(plan)
170    }
171
172    fn count(filter: &ListWorkflowsFilter) -> Result<Self, StoreError> {
173        Self::filtered("SELECT COUNT(*) FROM visibility", filter)
174    }
175
176    fn filtered(base_sql: &str, filter: &ListWorkflowsFilter) -> Result<Self, StoreError> {
177        let mut builder = FilterBuilder::default();
178        builder.add_filter(filter)?;
179        let (where_sql, params) = builder.finish();
180        let sql = if where_sql.is_empty() {
181            String::from(base_sql)
182        } else {
183            format!("{base_sql} WHERE {where_sql}")
184        };
185
186        Ok(Self { sql, params })
187    }
188
189    fn push_param(&mut self, value: Value) {
190        self.params.push(value);
191    }
192}
193
194#[derive(Default)]
195struct FilterBuilder {
196    clauses: Vec<String>,
197    params: Vec<Value>,
198}
199
200impl FilterBuilder {
201    fn add_filter(&mut self, filter: &ListWorkflowsFilter) -> Result<(), StoreError> {
202        if let Some(workflow_type) = &filter.workflow_type {
203            self.push_clause("workflow_type =", Value::Text(workflow_type.clone()));
204        }
205        if let Some(status) = filter.status {
206            self.push_clause("status =", Value::Text(encode_status(status)?));
207        }
208        if let Some(started_after) = filter.started_after {
209            self.push_clause(
210                "start_time >=",
211                Value::Text(encode_timestamp(started_after)),
212            );
213        }
214        if let Some(started_before) = filter.started_before {
215            self.push_clause(
216                "start_time <=",
217                Value::Text(encode_timestamp(started_before)),
218            );
219        }
220        if let Some(closed_after) = filter.closed_after {
221            self.push_clause("close_time >=", Value::Text(encode_timestamp(closed_after)));
222        }
223        if let Some(closed_before) = filter.closed_before {
224            self.push_clause(
225                "close_time <=",
226                Value::Text(encode_timestamp(closed_before)),
227            );
228        }
229        for predicate in &filter.search_attributes {
230            self.add_search_attribute_predicate(predicate)?;
231        }
232
233        Ok(())
234    }
235
236    fn push_clause(&mut self, lhs_and_operator: &str, value: Value) {
237        self.params.push(value);
238        self.clauses
239            .push(format!("{lhs_and_operator} ?{}", self.params.len()));
240    }
241
242    fn add_search_attribute_predicate(
243        &mut self,
244        predicate: &SearchAttributePredicate,
245    ) -> Result<(), StoreError> {
246        match predicate {
247            SearchAttributePredicate::Equals { name, value } => self.add_equals(name, value),
248            SearchAttributePredicate::GreaterThan { name, value } => {
249                self.add_ordered_comparison(name, value, ">")
250            }
251            SearchAttributePredicate::LessThan { name, value } => {
252                self.add_ordered_comparison(name, value, "<")
253            }
254            SearchAttributePredicate::Contains { name, keyword } => {
255                self.add_contains(name, keyword)
256            }
257        }
258    }
259
260    fn add_equals(&mut self, name: &str, value: &SearchAttributeValue) -> Result<(), StoreError> {
261        let type_path = search_attribute_path(name, "type")?;
262        let data_path = search_attribute_path(name, "data")?;
263        let type_param = self.push(Value::Text(type_name(value)));
264        let data_param = self.push(search_attribute_data_value(value)?);
265        let type_path_param = self.push(Value::Text(type_path));
266        let data_path_param = self.push(Value::Text(data_path));
267        self.clauses.push(format!(
268            "json_extract(search_attributes, ?{type_path_param}) = ?{type_param} AND json_extract(search_attributes, ?{data_path_param}) = ?{data_param}"
269        ));
270        Ok(())
271    }
272
273    fn add_ordered_comparison(
274        &mut self,
275        name: &str,
276        value: &SearchAttributeValue,
277        operator: &str,
278    ) -> Result<(), StoreError> {
279        if !is_ordered_value(value) {
280            self.clauses.push(String::from("0 = 1"));
281            return Ok(());
282        }
283
284        let type_path = search_attribute_path(name, "type")?;
285        let data_path = search_attribute_path(name, "data")?;
286        let type_param = self.push(Value::Text(type_name(value)));
287        let data_param = self.push(search_attribute_data_value(value)?);
288        let type_path_param = self.push(Value::Text(type_path));
289        let data_path_param = self.push(Value::Text(data_path));
290        self.clauses.push(format!(
291            "json_extract(search_attributes, ?{type_path_param}) = ?{type_param} AND json_extract(search_attributes, ?{data_path_param}) {operator} ?{data_param}"
292        ));
293        Ok(())
294    }
295
296    fn add_contains(&mut self, name: &str, keyword: &str) -> Result<(), StoreError> {
297        let type_path = search_attribute_path(name, "type")?;
298        let data_path = search_attribute_path(name, "data")?;
299        let type_param = self.push(Value::Text(String::from("KeywordList")));
300        let keyword_param = self.push(Value::Text(String::from(keyword)));
301        let type_path_param = self.push(Value::Text(type_path));
302        let data_path_param = self.push(Value::Text(data_path));
303        self.clauses.push(format!(
304            "json_extract(search_attributes, ?{type_path_param}) = ?{type_param} AND EXISTS (SELECT 1 FROM json_each(json_extract(search_attributes, ?{data_path_param})) WHERE value = ?{keyword_param})"
305        ));
306        Ok(())
307    }
308
309    fn push(&mut self, value: Value) -> usize {
310        self.params.push(value);
311        self.params.len()
312    }
313
314    fn finish(self) -> (String, Vec<Value>) {
315        (self.clauses.join(" AND "), self.params)
316    }
317}
318
319fn decode_summary(row: &libsql::Row) -> Result<VisibilityWorkflowSummary, StoreError> {
320    let workflow_id: String = row
321        .get(0)
322        .map_err(|error| crate::error::libsql_error(&error))?;
323    let run_id: String = row
324        .get(1)
325        .map_err(|error| crate::error::libsql_error(&error))?;
326    let workflow_type: String = row
327        .get(2)
328        .map_err(|error| crate::error::libsql_error(&error))?;
329    let status: String = row
330        .get(3)
331        .map_err(|error| crate::error::libsql_error(&error))?;
332    let start_time: String = row
333        .get(4)
334        .map_err(|error| crate::error::libsql_error(&error))?;
335    let close_time: Option<String> = row
336        .get(5)
337        .map_err(|error| crate::error::libsql_error(&error))?;
338    let failed_step: Option<String> = row
339        .get(6)
340        .map_err(|error| crate::error::libsql_error(&error))?;
341    let failure_reason: Option<String> = row
342        .get(7)
343        .map_err(|error| crate::error::libsql_error(&error))?;
344    let search_attributes: String = row
345        .get(8)
346        .map_err(|error| crate::error::libsql_error(&error))?;
347
348    Ok(VisibilityWorkflowSummary {
349        workflow_id: decode_workflow_id(&workflow_id)?,
350        run_id: decode_run_id(&run_id)?,
351        workflow_type,
352        status: decode_status(&status)?,
353        start_time: decode_timestamp(&start_time)?,
354        close_time: close_time.as_deref().map(decode_timestamp).transpose()?,
355        failed_step,
356        failure_reason,
357        search_attributes: serde_json::from_str::<HashMap<String, SearchAttributeValue>>(
358            &search_attributes,
359        )
360        .map_err(|error| crate::error::serde_json_error(&error))?,
361    })
362}
363
364fn encode_status(status: WorkflowStatus) -> Result<String, StoreError> {
365    serde_json::to_string(&status).map_err(|error| crate::error::serde_json_error(&error))
366}
367
368fn decode_status(value: &str) -> Result<WorkflowStatus, StoreError> {
369    serde_json::from_str(value).map_err(|error| crate::error::serde_json_error(&error))
370}
371
372fn encode_timestamp(timestamp: DateTime<Utc>) -> String {
373    timestamp.to_rfc3339_opts(SecondsFormat::Nanos, true)
374}
375
376fn decode_timestamp(value: &str) -> Result<DateTime<Utc>, StoreError> {
377    DateTime::parse_from_rfc3339(value)
378        .map(|date_time| date_time.with_timezone(&Utc))
379        .map_err(|error| StoreError::Serialization(error.to_string()))
380}
381
382fn decode_workflow_id(value: &str) -> Result<WorkflowId, StoreError> {
383    Uuid::parse_str(value)
384        .map(WorkflowId::new)
385        .map_err(|error| StoreError::Serialization(error.to_string()))
386}
387
388fn decode_run_id(value: &str) -> Result<RunId, StoreError> {
389    Uuid::parse_str(value)
390        .map(RunId::new)
391        .map_err(|error| StoreError::Serialization(error.to_string()))
392}
393
394fn search_attribute_path(name: &str, field: &str) -> Result<String, StoreError> {
395    let quoted_name =
396        serde_json::to_string(name).map_err(|error| crate::error::serde_json_error(&error))?;
397    Ok(format!("$.{quoted_name}.{field}"))
398}
399
400fn type_name(value: &SearchAttributeValue) -> String {
401    match value {
402        SearchAttributeValue::String(_) => String::from("String"),
403        SearchAttributeValue::Int(_) => String::from("Int"),
404        SearchAttributeValue::Float(_) => String::from("Float"),
405        SearchAttributeValue::Bool(_) => String::from("Bool"),
406        SearchAttributeValue::Datetime(_) => String::from("Datetime"),
407        SearchAttributeValue::KeywordList(_) => String::from("KeywordList"),
408    }
409}
410
411fn search_attribute_data_value(value: &SearchAttributeValue) -> Result<Value, StoreError> {
412    match value {
413        SearchAttributeValue::String(value) => Ok(Value::Text(value.clone())),
414        SearchAttributeValue::Int(value) => Ok(Value::Integer(*value)),
415        SearchAttributeValue::Float(value) => Ok(Value::Real(*value)),
416        SearchAttributeValue::Bool(value) => Ok(Value::Integer(i64::from(*value))),
417        SearchAttributeValue::Datetime(value) => {
418            Ok(Value::Text(serde_search_attribute_data(value)?))
419        }
420        SearchAttributeValue::KeywordList(values) => serde_json::to_string(values)
421            .map(Value::Text)
422            .map_err(|error| crate::error::serde_json_error(&error)),
423    }
424}
425
426fn serde_search_attribute_data(value: &DateTime<Utc>) -> Result<String, StoreError> {
427    let json = serde_json::to_value(SearchAttributeValue::Datetime(value.to_owned()))
428        .map_err(|error| crate::error::serde_json_error(&error))?;
429    json.get("data")
430        .and_then(serde_json::Value::as_str)
431        .map(String::from)
432        .ok_or_else(|| {
433            StoreError::Serialization(String::from(
434                "datetime search attribute encoded without string data",
435            ))
436        })
437}
438
439const fn is_ordered_value(value: &SearchAttributeValue) -> bool {
440    matches!(
441        value,
442        SearchAttributeValue::Int(_)
443            | SearchAttributeValue::Float(_)
444            | SearchAttributeValue::Datetime(_)
445    )
446}
447
448#[cfg(test)]
449mod tests {
450    use std::collections::HashMap;
451    use std::path::PathBuf;
452    use std::time::{SystemTime, UNIX_EPOCH};
453
454    use aion_core::{RunId, SearchAttributeValue, WorkflowId, WorkflowStatus};
455    use aion_store::StoreError;
456    use aion_store::visibility::{
457        ListWorkflowsFilter, SearchAttributePredicate, VisibilityRecord, VisibilityStore,
458        WorkflowSummary as VisibilityWorkflowSummary,
459    };
460    use chrono::{DateTime, TimeZone, Utc};
461
462    use super::{
463        count_workflows, decode_timestamp, encode_timestamp, list_workflows, record_visibility,
464    };
465    use crate::config::{LibSqlConfig, LibSqlMode};
466
467    #[tokio::test]
468    async fn record_visibility_inserts_and_updates_queryable_row() -> Result<(), StoreError> {
469        let conn = open_test_connection("upsert").await?;
470        let workflow_id = WorkflowId::new_v4();
471        let run_id = RunId::new_v4();
472        let mut record = visibility_record(
473            workflow_id.clone(),
474            run_id.clone(),
475            "orders",
476            WorkflowStatus::Running,
477            instant(2026, 6, 1, 9, 0, 0)?,
478            None,
479        );
480        record_visibility(&conn, record.clone()).await?;
481
482        record.status = WorkflowStatus::Completed;
483        record.close_time = Some(instant(2026, 6, 1, 10, 0, 0)?);
484        record.search_attributes.insert(
485            String::from("customer"),
486            SearchAttributeValue::String(String::from("cust-2")),
487        );
488        record_visibility(&conn, record.clone()).await?;
489
490        let summaries = list_workflows(&conn, ListWorkflowsFilter::default()).await?;
491        assert_eq!(summaries.len(), 1);
492        assert_eq!(summaries[0], record.into());
493
494        Ok(())
495    }
496
497    #[tokio::test]
498    async fn list_filters_standard_fields_and_time_ranges() -> Result<(), StoreError> {
499        let conn = open_test_connection("standard-filters").await?;
500        let first = seed_record(
501            &conn,
502            "orders",
503            WorkflowStatus::Completed,
504            instant(2026, 6, 1, 9, 0, 0)?,
505            Some(instant(2026, 6, 1, 10, 0, 0)?),
506            "cust-1",
507            2,
508        )
509        .await?;
510        let second = seed_record(
511            &conn,
512            "billing",
513            WorkflowStatus::Running,
514            instant(2026, 6, 2, 9, 0, 0)?,
515            None,
516            "cust-2",
517            5,
518        )
519        .await?;
520
521        assert_eq!(
522            ids(list_workflows(
523                &conn,
524                ListWorkflowsFilter {
525                    workflow_type: Some(String::from("orders")),
526                    ..ListWorkflowsFilter::default()
527                },
528            )
529            .await?),
530            vec![first.workflow_id.clone()]
531        );
532        assert_eq!(
533            ids(list_workflows(
534                &conn,
535                ListWorkflowsFilter {
536                    status: Some(WorkflowStatus::Running),
537                    ..ListWorkflowsFilter::default()
538                },
539            )
540            .await?),
541            vec![second.workflow_id.clone()]
542        );
543        assert_eq!(
544            ids(list_workflows(
545                &conn,
546                ListWorkflowsFilter {
547                    started_after: Some(instant(2026, 6, 2, 0, 0, 0)?),
548                    started_before: Some(instant(2026, 6, 2, 23, 59, 59)?),
549                    ..ListWorkflowsFilter::default()
550                },
551            )
552            .await?),
553            vec![second.workflow_id.clone()]
554        );
555        assert_eq!(
556            ids(list_workflows(
557                &conn,
558                ListWorkflowsFilter {
559                    closed_after: Some(instant(2026, 6, 1, 9, 30, 0)?),
560                    closed_before: Some(instant(2026, 6, 1, 10, 30, 0)?),
561                    ..ListWorkflowsFilter::default()
562                },
563            )
564            .await?),
565            vec![first.workflow_id.clone()]
566        );
567
568        Ok(())
569    }
570
571    #[tokio::test]
572    async fn list_filters_custom_search_attributes() -> Result<(), StoreError> {
573        let conn = open_test_connection("custom-filters").await?;
574        let first = seed_record(
575            &conn,
576            "orders",
577            WorkflowStatus::Completed,
578            instant(2026, 6, 1, 9, 0, 0)?,
579            Some(instant(2026, 6, 1, 10, 0, 0)?),
580            "cust-1",
581            2,
582        )
583        .await?;
584        let second = seed_record(
585            &conn,
586            "orders",
587            WorkflowStatus::Completed,
588            instant(2026, 6, 2, 9, 0, 0)?,
589            Some(instant(2026, 6, 2, 10, 0, 0)?),
590            "cust-2",
591            5,
592        )
593        .await?;
594
595        assert_eq!(
596            ids(list_workflows(
597                &conn,
598                custom_filter(SearchAttributePredicate::Equals {
599                    name: String::from("customer"),
600                    value: SearchAttributeValue::String(String::from("cust-1")),
601                })
602            )
603            .await?),
604            vec![first.workflow_id.clone()]
605        );
606        assert_eq!(
607            ids(list_workflows(
608                &conn,
609                custom_filter(SearchAttributePredicate::GreaterThan {
610                    name: String::from("attempts"),
611                    value: SearchAttributeValue::Int(3),
612                })
613            )
614            .await?),
615            vec![second.workflow_id.clone()]
616        );
617        assert_eq!(
618            ids(list_workflows(
619                &conn,
620                custom_filter(SearchAttributePredicate::LessThan {
621                    name: String::from("attempts"),
622                    value: SearchAttributeValue::Int(3),
623                })
624            )
625            .await?),
626            vec![first.workflow_id.clone()]
627        );
628        assert_eq!(
629            ids(list_workflows(
630                &conn,
631                custom_filter(SearchAttributePredicate::Contains {
632                    name: String::from("tags"),
633                    keyword: String::from("west"),
634                })
635            )
636            .await?),
637            vec![second.workflow_id.clone(), first.workflow_id.clone()]
638        );
639
640        Ok(())
641    }
642
643    #[tokio::test]
644    async fn list_empty_filter_orders_by_start_time_desc_and_paginates() -> Result<(), StoreError> {
645        let conn = open_test_connection("pagination").await?;
646        let oldest = seed_record(
647            &conn,
648            "orders",
649            WorkflowStatus::Running,
650            instant(2026, 6, 1, 9, 0, 0)?,
651            None,
652            "cust-1",
653            1,
654        )
655        .await?;
656        let middle = seed_record(
657            &conn,
658            "orders",
659            WorkflowStatus::Running,
660            instant(2026, 6, 2, 9, 0, 0)?,
661            None,
662            "cust-2",
663            2,
664        )
665        .await?;
666        let newest = seed_record(
667            &conn,
668            "orders",
669            WorkflowStatus::Running,
670            instant(2026, 6, 3, 9, 0, 0)?,
671            None,
672            "cust-3",
673            3,
674        )
675        .await?;
676
677        assert_eq!(
678            ids(list_workflows(&conn, ListWorkflowsFilter::default()).await?),
679            vec![
680                newest.workflow_id,
681                middle.workflow_id.clone(),
682                oldest.workflow_id
683            ]
684        );
685        assert_eq!(
686            ids(list_workflows(
687                &conn,
688                ListWorkflowsFilter {
689                    limit: Some(1),
690                    offset: Some(1),
691                    ..ListWorkflowsFilter::default()
692                },
693            )
694            .await?),
695            vec![middle.workflow_id]
696        );
697
698        Ok(())
699    }
700
701    #[tokio::test]
702    async fn count_workflows_counts_total_and_filtered_matches() -> Result<(), StoreError> {
703        let conn = open_test_connection("count").await?;
704        seed_record(
705            &conn,
706            "orders",
707            WorkflowStatus::Completed,
708            instant(2026, 6, 1, 9, 0, 0)?,
709            Some(instant(2026, 6, 1, 10, 0, 0)?),
710            "cust-1",
711            2,
712        )
713        .await?;
714        seed_record(
715            &conn,
716            "orders",
717            WorkflowStatus::Running,
718            instant(2026, 6, 2, 9, 0, 0)?,
719            None,
720            "cust-2",
721            5,
722        )
723        .await?;
724
725        assert_eq!(
726            count_workflows(&conn, ListWorkflowsFilter::default()).await?,
727            2
728        );
729        let filter = ListWorkflowsFilter {
730            status: Some(WorkflowStatus::Completed),
731            ..ListWorkflowsFilter::default()
732        };
733        assert_eq!(count_workflows(&conn, filter.clone()).await?, 1);
734        assert_eq!(
735            usize::try_from(count_workflows(&conn, filter.clone()).await?).unwrap_or(usize::MAX),
736            list_workflows(&conn, filter).await?.len()
737        );
738
739        Ok(())
740    }
741
742    #[tokio::test]
743    async fn libsql_store_implements_visibility_store_trait() -> Result<(), StoreError> {
744        let store = crate::store::LibSqlStore::open(unique_temp_path("trait")).await?;
745        let store: std::sync::Arc<dyn VisibilityStore> = std::sync::Arc::new(store);
746
747        assert_eq!(std::sync::Arc::strong_count(&store), 1);
748        Ok(())
749    }
750
751    #[test]
752    fn timestamp_encoding_round_trips_losslessly() -> Result<(), StoreError> {
753        let timestamp = instant(2026, 6, 3, 2, 30, 0)?;
754
755        assert_eq!(decode_timestamp(&encode_timestamp(timestamp))?, timestamp);
756        Ok(())
757    }
758
759    async fn seed_record(
760        conn: &libsql::Connection,
761        workflow_type: &str,
762        status: WorkflowStatus,
763        start_time: DateTime<Utc>,
764        close_time: Option<DateTime<Utc>>,
765        customer: &str,
766        attempts: i64,
767    ) -> Result<VisibilityRecord, StoreError> {
768        let record = visibility_record(
769            WorkflowId::new_v4(),
770            RunId::new_v4(),
771            workflow_type,
772            status,
773            start_time,
774            close_time,
775        )
776        .with_customer(customer, attempts);
777        record_visibility(conn, record.clone()).await?;
778        Ok(record)
779    }
780
781    trait RecordBuilder {
782        fn with_customer(self, customer: &str, attempts: i64) -> Self;
783    }
784
785    impl RecordBuilder for VisibilityRecord {
786        fn with_customer(mut self, customer: &str, attempts: i64) -> Self {
787            self.search_attributes.insert(
788                String::from("customer"),
789                SearchAttributeValue::String(String::from(customer)),
790            );
791            self.search_attributes.insert(
792                String::from("attempts"),
793                SearchAttributeValue::Int(attempts),
794            );
795            self.search_attributes.insert(
796                String::from("tags"),
797                SearchAttributeValue::KeywordList(vec![String::from("vip"), String::from("west")]),
798            );
799            self
800        }
801    }
802
803    fn visibility_record(
804        workflow_id: WorkflowId,
805        run_id: RunId,
806        workflow_type: &str,
807        status: WorkflowStatus,
808        start_time: DateTime<Utc>,
809        close_time: Option<DateTime<Utc>>,
810    ) -> VisibilityRecord {
811        VisibilityRecord {
812            workflow_id,
813            run_id,
814            workflow_type: String::from(workflow_type),
815            status,
816            start_time,
817            close_time,
818            failed_step: None,
819            failure_reason: None,
820            search_attributes: HashMap::new(),
821        }
822    }
823
824    #[tokio::test]
825    async fn record_visibility_round_trips_failed_step_and_reason() -> Result<(), StoreError> {
826        let conn = open_test_connection("failed-fields").await?;
827        let mut record = visibility_record(
828            WorkflowId::new_v4(),
829            RunId::new_v4(),
830            "stacked_dev",
831            WorkflowStatus::Failed,
832            instant(2026, 6, 1, 9, 0, 0)?,
833            Some(instant(2026, 6, 1, 10, 0, 0)?),
834        );
835        record.failed_step = Some(String::from("dev_review"));
836        record.failure_reason = Some(String::from("provider error: rate limited"));
837        record_visibility(&conn, record.clone()).await?;
838
839        let summaries = list_workflows(&conn, ListWorkflowsFilter::default()).await?;
840        assert_eq!(summaries.len(), 1);
841        assert_eq!(summaries[0].failed_step.as_deref(), Some("dev_review"));
842        assert_eq!(
843            summaries[0].failure_reason.as_deref(),
844            Some("provider error: rate limited")
845        );
846        assert_eq!(summaries[0], record.into());
847        Ok(())
848    }
849
850    fn custom_filter(predicate: SearchAttributePredicate) -> ListWorkflowsFilter {
851        ListWorkflowsFilter {
852            search_attributes: vec![predicate],
853            ..ListWorkflowsFilter::default()
854        }
855    }
856
857    fn ids(summaries: Vec<VisibilityWorkflowSummary>) -> Vec<WorkflowId> {
858        summaries
859            .into_iter()
860            .map(|summary| summary.workflow_id)
861            .collect()
862    }
863
864    async fn open_test_connection(name: &str) -> Result<libsql::Connection, StoreError> {
865        let config = LibSqlConfig {
866            mode: LibSqlMode::Embedded {
867                path: unique_temp_path(name),
868            },
869            journal_mode: None,
870            synchronous: None,
871            sync_interval_seconds: None,
872        };
873        let conn = crate::connection::open_connection(&config)
874            .await?
875            .connection;
876        crate::schema::ensure_schema(&conn).await?;
877        Ok(conn)
878    }
879
880    fn instant(
881        year: i32,
882        month: u32,
883        day: u32,
884        hour: u32,
885        minute: u32,
886        second: u32,
887    ) -> Result<DateTime<Utc>, StoreError> {
888        Utc.with_ymd_and_hms(year, month, day, hour, minute, second)
889            .single()
890            .ok_or_else(|| StoreError::Serialization(String::from("invalid test instant")))
891    }
892
893    fn unique_temp_path(name: &str) -> PathBuf {
894        let nanos = SystemTime::now()
895            .duration_since(UNIX_EPOCH)
896            .map_or(0, |duration| duration.as_nanos());
897        std::env::temp_dir().join(format!(
898            "aion-store-libsql-visibility-{name}-{}-{nanos}.db",
899            std::process::id()
900        ))
901    }
902}