Add error_code and search_key filters to workflow runs API (#SKY-7884) (#4694)

Co-authored-by: Suchintan Singh <suchintan@skyvern.com>
This commit is contained in:
Suchintan
2026-02-11 00:42:11 -05:00
committed by GitHub
parent 9f7311ccd0
commit b1758dd3b5
3 changed files with 128 additions and 37 deletions

View File

@@ -3,7 +3,24 @@ from datetime import datetime, timedelta
from typing import Any, List, Literal, Sequence, overload
import structlog
from sqlalchemy import Text, and_, asc, case, delete, distinct, exists, func, or_, pool, select, tuple_, update
from sqlalchemy import (
Text,
and_,
asc,
case,
cast,
delete,
distinct,
exists,
func,
literal,
or_,
pool,
select,
tuple_,
update,
)
from sqlalchemy.dialects.postgresql import JSONB
from sqlalchemy.exc import (
SQLAlchemyError,
)
@@ -3011,6 +3028,57 @@ class AgentDB(BaseAlchemyDB):
LOG.error("SQLAlchemyError", exc_info=True)
raise
@staticmethod
def _apply_search_key_filter(query, search_key: str | None): # type: ignore[no-untyped-def]
if not search_key:
return query
key_like = f"%{search_key}%"
# Match workflow_run_id directly
id_matches = WorkflowRunModel.workflow_run_id.ilike(key_like)
# Match parameter key or description (only for non-deleted parameter definitions)
# Use EXISTS to avoid duplicate rows and to keep pagination correct
param_key_desc_exists = exists(
select(1)
.select_from(WorkflowRunParameterModel)
.join(
WorkflowParameterModel,
WorkflowParameterModel.workflow_parameter_id == WorkflowRunParameterModel.workflow_parameter_id,
)
.where(WorkflowRunParameterModel.workflow_run_id == WorkflowRunModel.workflow_run_id)
.where(WorkflowParameterModel.deleted_at.is_(None))
.where(
or_(
WorkflowParameterModel.key.ilike(key_like),
WorkflowParameterModel.description.ilike(key_like),
)
)
)
# Match run parameter value directly (searches all values regardless of parameter definition status)
param_value_exists = exists(
select(1)
.select_from(WorkflowRunParameterModel)
.where(WorkflowRunParameterModel.workflow_run_id == WorkflowRunModel.workflow_run_id)
.where(WorkflowRunParameterModel.value.ilike(key_like))
)
# Match extra HTTP headers (cast JSON to text for search, skip NULLs)
extra_headers_match = and_(
WorkflowRunModel.extra_http_headers.isnot(None),
func.cast(WorkflowRunModel.extra_http_headers, Text()).ilike(key_like),
)
return query.where(or_(id_matches, param_key_desc_exists, param_value_exists, extra_headers_match))
@staticmethod
def _apply_error_code_filter(query, error_code: str | None): # type: ignore[no-untyped-def]
if not error_code:
return query
error_code_exists = exists(
select(1)
.select_from(TaskModel)
.where(TaskModel.workflow_run_id == WorkflowRunModel.workflow_run_id)
.where(cast(TaskModel.errors, JSONB).contains(literal([{"error_code": error_code}], type_=JSONB)))
)
return query.where(error_code_exists)
async def get_workflow_runs(
self,
organization_id: str,
@@ -3018,6 +3086,8 @@ class AgentDB(BaseAlchemyDB):
page_size: int = 10,
status: list[WorkflowRunStatus] | None = None,
ordering: tuple[str, str] | None = None,
search_key: str | None = None,
error_code: str | None = None,
) -> list[WorkflowRun]:
try:
async with self.Session() as session:
@@ -3030,6 +3100,9 @@ class AgentDB(BaseAlchemyDB):
.filter(WorkflowRunModel.parent_workflow_run_id.is_(None))
)
query = self._apply_search_key_filter(query, search_key)
query = self._apply_error_code_filter(query, error_code)
if status:
query = query.filter(WorkflowRunModel.status.in_(status))
@@ -3092,6 +3165,7 @@ class AgentDB(BaseAlchemyDB):
page_size: int = 10,
status: list[WorkflowRunStatus] | None = None,
search_key: str | None = None,
error_code: str | None = None,
) -> list[WorkflowRun]:
"""
Get runs for a workflow, with optional `search_key` on run ID, parameter key/description/value,
@@ -3106,42 +3180,8 @@ class AgentDB(BaseAlchemyDB):
.filter(WorkflowRunModel.workflow_permanent_id == workflow_permanent_id)
.filter(WorkflowRunModel.organization_id == organization_id)
)
if search_key:
key_like = f"%{search_key}%"
# Match workflow_run_id directly
id_matches = WorkflowRunModel.workflow_run_id.ilike(key_like)
# Match parameter key or description (only for non-deleted parameter definitions)
# Use EXISTS to avoid duplicate rows and to keep pagination correct
param_key_desc_exists = exists(
select(1)
.select_from(WorkflowRunParameterModel)
.join(
WorkflowParameterModel,
WorkflowParameterModel.workflow_parameter_id
== WorkflowRunParameterModel.workflow_parameter_id,
)
.where(WorkflowRunParameterModel.workflow_run_id == WorkflowRunModel.workflow_run_id)
.where(WorkflowParameterModel.deleted_at.is_(None))
.where(
or_(
WorkflowParameterModel.key.ilike(key_like),
WorkflowParameterModel.description.ilike(key_like),
)
)
)
# Match run parameter value directly (searches all values regardless of parameter definition status)
param_value_exists = exists(
select(1)
.select_from(WorkflowRunParameterModel)
.where(WorkflowRunParameterModel.workflow_run_id == WorkflowRunModel.workflow_run_id)
.where(WorkflowRunParameterModel.value.ilike(key_like))
)
# Match extra HTTP headers (cast JSON to text for search, skip NULLs)
extra_headers_match = and_(
WorkflowRunModel.extra_http_headers.isnot(None),
func.cast(WorkflowRunModel.extra_http_headers, Text()).ilike(key_like),
)
query = query.where(or_(id_matches, param_key_desc_exists, param_value_exists, extra_headers_match))
query = self._apply_search_key_filter(query, search_key)
query = self._apply_error_code_filter(query, error_code)
if status:
query = query.filter(WorkflowRunModel.status.in_(status))
query = query.order_by(WorkflowRunModel.created_at.desc()).limit(page_size).offset(db_page * page_size)