diff --git a/olymp/events.py b/olymp/events.py index 053fd6f..f1da745 100644 --- a/olymp/events.py +++ b/olymp/events.py @@ -26,20 +26,56 @@ _LIST_EVENTS_QUERY = """ SELECT * FROM events - WHERE (? IS NULL OR event_type = ?) + ORDER BY event_id ASC + LIMIT ? +""" +_LIST_EVENTS_AFTER_QUERY = """ + SELECT * + FROM events + WHERE event_id > ? + AND (? IS NULL OR event_type = ?) AND (? IS NULL OR run_id = ?) AND (? IS NULL OR plan_id = ?) AND (? IS NULL OR node_id = ?) ORDER BY event_id ASC LIMIT ? """ -_LIST_EVENTS_AFTER_QUERY = """ +_LIST_EVENTS_BY_RUN_QUERY = """ SELECT * FROM events - WHERE event_id > ? + WHERE run_id = ? + AND (? IS NULL OR event_type = ?) + AND (? IS NULL OR plan_id = ?) + AND (? IS NULL OR node_id = ?) + ORDER BY event_id ASC + LIMIT ? +""" +_LIST_EVENTS_BY_PLAN_QUERY = """ + SELECT * + FROM events + WHERE plan_id = ? + AND (? IS NULL OR event_type = ?) + AND (? IS NULL OR run_id = ?) + AND (? IS NULL OR node_id = ?) + ORDER BY event_id ASC + LIMIT ? +""" +_LIST_EVENTS_BY_NODE_QUERY = """ + SELECT * + FROM events + WHERE node_id = ? AND (? IS NULL OR event_type = ?) AND (? IS NULL OR run_id = ?) AND (? IS NULL OR plan_id = ?) + ORDER BY event_id ASC + LIMIT ? +""" +_LIST_EVENTS_BY_TYPE_QUERY = """ + SELECT * + FROM events + WHERE event_type = ? + AND (? IS NULL OR run_id = ?) + AND (? IS NULL OR plan_id = ?) AND (? IS NULL OR node_id = ?) ORDER BY event_id ASC LIMIT ? @@ -149,23 +185,71 @@ def list( checked_run = _nullable_text(run_id, "run_id") checked_plan = _nullable_text(plan_id, "plan_id") checked_node = _nullable_text(node_id, "node_id") - filter_values: tuple[object, ...] = ( - checked_type, - checked_type, - checked_run, - checked_run, - checked_plan, - checked_plan, - checked_node, - checked_node, - checked_limit, - ) - if checked_after is None: - query = _LIST_EVENTS_QUERY - values = filter_values - else: + if checked_after is not None: query = _LIST_EVENTS_AFTER_QUERY - values = (checked_after, *filter_values) + values: tuple[object, ...] = ( + checked_after, + checked_type, + checked_type, + checked_run, + checked_run, + checked_plan, + checked_plan, + checked_node, + checked_node, + checked_limit, + ) + elif checked_run is not None: + query = _LIST_EVENTS_BY_RUN_QUERY + values = ( + checked_run, + checked_type, + checked_type, + checked_plan, + checked_plan, + checked_node, + checked_node, + checked_limit, + ) + elif checked_plan is not None: + query = _LIST_EVENTS_BY_PLAN_QUERY + values = ( + checked_plan, + checked_type, + checked_type, + checked_run, + checked_run, + checked_node, + checked_node, + checked_limit, + ) + elif checked_node is not None: + query = _LIST_EVENTS_BY_NODE_QUERY + values = ( + checked_node, + checked_type, + checked_type, + checked_run, + checked_run, + checked_plan, + checked_plan, + checked_limit, + ) + elif checked_type is not None: + query = _LIST_EVENTS_BY_TYPE_QUERY + values = ( + checked_type, + checked_run, + checked_run, + checked_plan, + checked_plan, + checked_node, + checked_node, + checked_limit, + ) + else: + query = _LIST_EVENTS_QUERY + values = (checked_limit,) with closing(self._connect()) as db: rows = db.execute(query, values).fetchall() return [_event_record(row) for row in rows] diff --git a/tests/test_events.py b/tests/test_events.py index f4742e4..d988fff 100644 --- a/tests/test_events.py +++ b/tests/test_events.py @@ -93,6 +93,10 @@ def test_event_store_filters_and_paginates(self) -> None: "EXPLAIN QUERY PLAN " + events_module._LIST_EVENTS_AFTER_QUERY, (1, None, None, None, None, None, None, None, None, 100), ).fetchall() + run_query_plan = db.execute( + "EXPLAIN QUERY PLAN " + events_module._LIST_EVENTS_BY_RUN_QUERY, + ("run-a", None, None, None, None, None, None, 100), + ).fetchall() self.assertEqual([event["event_id"] for event in first_page], matching_ids[:1]) self.assertEqual([event["event_id"] for event in second_page], matching_ids[1:]) @@ -101,6 +105,13 @@ def test_event_store_filters_and_paginates(self) -> None: any("SEARCH events USING INTEGER PRIMARY KEY" in str(row[3]) for row in query_plan), query_plan, ) + self.assertTrue( + any( + "SEARCH events USING INDEX idx_events_run_id" in str(row[3]) + for row in run_query_plan + ), + run_query_plan, + ) def test_subscriber_failure_is_isolated_and_secret_safe(self) -> None: with tempfile.TemporaryDirectory() as temp: