-
-
Notifications
You must be signed in to change notification settings - Fork 2.1k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Logs: get_query_results() and describe_queries() (#6730)
- Loading branch information
Showing
10 changed files
with
491 additions
and
45 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,90 @@ | ||
from typing import Any, Dict, List | ||
from typing import TYPE_CHECKING | ||
|
||
if TYPE_CHECKING: | ||
from ..models import LogGroup, LogEvent, LogStream | ||
|
||
from .query_parser import parse_query, ParsedQuery | ||
|
||
|
||
class ParsedEvent: | ||
def __init__( | ||
self, | ||
event: "LogEvent", | ||
query: ParsedQuery, | ||
log_stream: "LogStream", | ||
log_group: "LogGroup", | ||
): | ||
self.event = event | ||
self.query = query | ||
self.log_stream = log_stream | ||
self.log_group = log_group | ||
self.fields = self._create_fields() | ||
|
||
def _create_fields(self) -> Dict[str, Any]: | ||
fields: Dict[str, Any] = {"@ptr": self.event.event_id} | ||
if "@timestamp" in self.query.fields: | ||
fields["@timestamp"] = self.event.timestamp | ||
if "@message" in self.query.fields: | ||
fields["@message"] = self.event.message | ||
if "@logStream" in self.query.fields: | ||
fields["@logStream"] = self.log_stream.log_stream_name # type: ignore[has-type] | ||
if "@log" in self.query.fields: | ||
fields["@log"] = self.log_group.name | ||
return fields | ||
|
||
def __eq__(self, other: "ParsedEvent") -> bool: # type: ignore[override] | ||
return self.event.timestamp == other.event.timestamp | ||
|
||
def __lt__(self, other: "ParsedEvent") -> bool: | ||
return self.event.timestamp < other.event.timestamp | ||
|
||
def __le__(self, other: "ParsedEvent") -> bool: | ||
return self.event.timestamp <= other.event.timestamp | ||
|
||
def __gt__(self, other: "ParsedEvent") -> bool: | ||
return self.event.timestamp > other.event.timestamp | ||
|
||
def __ge__(self, other: "ParsedEvent") -> bool: | ||
return self.event.timestamp >= other.event.timestamp | ||
|
||
|
||
def execute_query( | ||
log_groups: List["LogGroup"], query: str, start_time: int, end_time: int | ||
) -> List[Dict[str, str]]: | ||
parsed = parse_query(query) | ||
all_events = _create_parsed_events(log_groups, parsed, start_time, end_time) | ||
sorted_events = sorted(all_events, reverse=parsed.sort_reversed()) | ||
sorted_fields = [event.fields for event in sorted_events] | ||
if parsed.limit: | ||
return sorted_fields[0 : parsed.limit] | ||
return sorted_fields | ||
|
||
|
||
def _create_parsed_events( | ||
log_groups: List["LogGroup"], query: ParsedQuery, start_time: int, end_time: int | ||
) -> List["ParsedEvent"]: | ||
def filter_func(event: "LogEvent") -> bool: | ||
# Start/End time is in epoch seconds | ||
# Event timestamp is in epoch milliseconds | ||
if start_time and event.timestamp < (start_time * 1000): | ||
return False | ||
|
||
if end_time and event.timestamp > (end_time * 1000): | ||
return False | ||
|
||
return True | ||
|
||
events: List["ParsedEvent"] = [] | ||
for group in log_groups: | ||
for stream in group.streams.values(): | ||
events.extend( | ||
[ | ||
ParsedEvent( | ||
event=event, query=query, log_stream=stream, log_group=group | ||
) | ||
for event in filter(filter_func, stream.events) | ||
] | ||
) | ||
|
||
return events |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,74 @@ | ||
from typing import List, Optional, Tuple | ||
|
||
from moto.utilities.tokenizer import GenericTokenizer | ||
|
||
|
||
class ParsedQuery: | ||
def __init__(self) -> None: | ||
self.limit: Optional[int] = None | ||
self.fields: List[str] = [] | ||
self.sort: List[Tuple[str, str]] = [] | ||
|
||
def sort_reversed(self) -> bool: | ||
# Descending is the default | ||
if self.sort: | ||
# sort_reversed is True if we want to sort in ascending order | ||
return self.sort[-1][-1] == "asc" | ||
return False | ||
|
||
|
||
def parse_query(query: str) -> ParsedQuery: | ||
tokenizer = GenericTokenizer(query) | ||
state = "COMMAND" | ||
characters = "" | ||
parsed_query = ParsedQuery() | ||
|
||
for char in tokenizer: | ||
if char.isspace(): | ||
if state == "SORT": | ||
parsed_query.sort.append((characters, "desc")) | ||
characters = "" | ||
state = "SORT_ORDER" | ||
if state == "COMMAND": | ||
if characters.lower() in ["fields", "limit", "sort"]: | ||
state = characters.upper() | ||
else: | ||
# Unknown/Unsupported command | ||
pass | ||
characters = "" | ||
tokenizer.skip_white_space() | ||
continue | ||
|
||
if char == "|": | ||
if state == "FIELDS": | ||
parsed_query.fields.append(characters) | ||
characters = "" | ||
if state == "LIMIT": | ||
parsed_query.limit = int(characters) | ||
characters = "" | ||
if state == "SORT_ORDER": | ||
if characters != "": | ||
parsed_query.sort[-1] = (parsed_query.sort[-1][0], characters) | ||
characters = "" | ||
state = "COMMAND" | ||
tokenizer.skip_white_space() | ||
continue | ||
|
||
if char == ",": | ||
if state == "FIELDS": | ||
parsed_query.fields.append(characters) | ||
characters = "" | ||
continue | ||
|
||
characters += char | ||
|
||
if state == "FIELDS": | ||
parsed_query.fields.append(characters) | ||
if state == "LIMIT": | ||
parsed_query.limit = int(characters) | ||
if state == "SORT": | ||
parsed_query.sort.append((characters, "desc")) | ||
if state == "SORT_ORDER": | ||
parsed_query.sort[-1] = (parsed_query.sort[-1][0], characters) | ||
|
||
return parsed_query |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.