1use 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
33pub(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
73pub(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
100pub(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}