Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
122 changes: 103 additions & 19 deletions olymp/events.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 ?
Expand Down Expand Up @@ -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]
Expand Down
11 changes: 11 additions & 0 deletions tests/test_events.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:])
Expand All @@ -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:
Expand Down
Loading