|
| 1 | +from typing import Union, List, Dict |
| 2 | + |
| 3 | +from lumigo_tracer.parsing_utils import recursive_get_key, safe_get |
| 4 | +from lumigo_tracer.lumigo_utils import ( |
| 5 | + lumigo_safe_execute, |
| 6 | + Configuration, |
| 7 | + STEP_FUNCTION_UID_KEY, |
| 8 | + LUMIGO_EVENT_KEY, |
| 9 | + md5hash, |
| 10 | +) |
| 11 | + |
| 12 | +TRIGGER_CREATION_TIME_KEY = "approxEventCreationTime" |
| 13 | +MESSAGE_ID_KEY = "messageId" |
| 14 | +MESSAGE_IDS_KEY = "messageIds" |
| 15 | + |
| 16 | + |
| 17 | +def parse_triggered_by(event: dict): |
| 18 | + """ |
| 19 | + This function parses the event and build the dictionary that describes the given event. |
| 20 | +
|
| 21 | + The current possible values are: |
| 22 | + * {triggeredBy: unknown} |
| 23 | + * {triggeredBy: apigw, api: <host>, resource: <>, httpMethod: <>, stage: <>, identity: <>, referer: <>} |
| 24 | + """ |
| 25 | + with lumigo_safe_execute("triggered by"): |
| 26 | + if not isinstance(event, dict): |
| 27 | + if _is_step_function(event): |
| 28 | + return _parse_step_function(event) |
| 29 | + return None |
| 30 | + if _is_supported_http_method(event): |
| 31 | + return _parse_http_method(event) |
| 32 | + elif _is_supported_sns(event): |
| 33 | + return _parse_sns(event) |
| 34 | + elif _is_supported_streams(event): |
| 35 | + return _parse_streams(event) |
| 36 | + elif _is_supported_cw(event): |
| 37 | + return _parse_cw(event) |
| 38 | + elif _is_step_function(event): |
| 39 | + return _parse_step_function(event) |
| 40 | + |
| 41 | + return _parse_unknown(event) |
| 42 | + |
| 43 | + |
| 44 | +def _parse_unknown(event: dict): |
| 45 | + result = {"triggeredBy": "unknown"} |
| 46 | + return result |
| 47 | + |
| 48 | + |
| 49 | +def _is_step_function(event: Union[List, Dict]): |
| 50 | + return ( |
| 51 | + Configuration.is_step_function |
| 52 | + and isinstance(event, (list, dict)) # noqa |
| 53 | + and STEP_FUNCTION_UID_KEY in recursive_get_key(event, LUMIGO_EVENT_KEY, default={}) # noqa |
| 54 | + ) |
| 55 | + |
| 56 | + |
| 57 | +def _parse_step_function(event: dict): |
| 58 | + result = { |
| 59 | + "triggeredBy": "stepFunction", |
| 60 | + "messageId": recursive_get_key(event, LUMIGO_EVENT_KEY)[STEP_FUNCTION_UID_KEY], |
| 61 | + } |
| 62 | + return result |
| 63 | + |
| 64 | + |
| 65 | +def _is_supported_http_method(event: dict): |
| 66 | + return ( |
| 67 | + "httpMethod" in event # noqa |
| 68 | + and "headers" in event # noqa |
| 69 | + and "requestContext" in event # noqa |
| 70 | + and event.get("requestContext", {}).get("elb") is None # noqa |
| 71 | + ) or ( # noqa |
| 72 | + event.get("version", "") == "2.0" and "headers" in event # noqa |
| 73 | + ) # noqa # noqa |
| 74 | + |
| 75 | + |
| 76 | +def _parse_http_method(event: dict): |
| 77 | + version = event.get("version") |
| 78 | + if version and version.startswith("2.0"): |
| 79 | + return _parse_http_method_v2(event) |
| 80 | + return _parse_http_method_v1(event) |
| 81 | + |
| 82 | + |
| 83 | +def _parse_http_method_v1(event: dict): |
| 84 | + result = { |
| 85 | + "triggeredBy": "apigw", |
| 86 | + "httpMethod": event.get("httpMethod", ""), |
| 87 | + "resource": event.get("resource", ""), |
| 88 | + "messageId": event.get("requestContext", {}).get("requestId", ""), |
| 89 | + } |
| 90 | + if isinstance(event.get("headers"), dict): |
| 91 | + result["api"] = event["headers"].get("Host", "unknown.unknown.unknown") |
| 92 | + if isinstance(event.get("requestContext"), dict): |
| 93 | + result["stage"] = event["requestContext"].get("stage", "unknown") |
| 94 | + return result |
| 95 | + |
| 96 | + |
| 97 | +def _parse_http_method_v2(event: dict): |
| 98 | + result = { |
| 99 | + "triggeredBy": "apigw", |
| 100 | + "httpMethod": event.get("requestContext", {}).get("http", {}).get("method"), |
| 101 | + "resource": event.get("requestContext", {}).get("http", {}).get("path"), |
| 102 | + "messageId": event.get("requestContext", {}).get("requestId", ""), |
| 103 | + "api": event.get("requestContext", {}).get("domainName", ""), |
| 104 | + "stage": event.get("requestContext", {}).get("stage", "unknown"), |
| 105 | + } |
| 106 | + return result |
| 107 | + |
| 108 | + |
| 109 | +def _is_supported_sns(event: dict): |
| 110 | + return event.get("Records", [{}])[0].get("EventSource") == "aws:sns" |
| 111 | + |
| 112 | + |
| 113 | +def _parse_sns(event: dict): |
| 114 | + return { |
| 115 | + "triggeredBy": "sns", |
| 116 | + "arn": event["Records"][0]["Sns"]["TopicArn"], |
| 117 | + "messageId": event["Records"][0]["Sns"].get("MessageId"), |
| 118 | + } |
| 119 | + |
| 120 | + |
| 121 | +def _is_supported_cw(event: dict): |
| 122 | + return event.get("detail-type") == "Scheduled Event" and "source" in event and "time" in event |
| 123 | + |
| 124 | + |
| 125 | +def _parse_cw(event: dict): |
| 126 | + resource = event.get("resources", ["/unknown"])[0].split("/")[1] |
| 127 | + return { |
| 128 | + "triggeredBy": "cloudwatch", |
| 129 | + "resource": resource, |
| 130 | + "region": event.get("region"), |
| 131 | + "detailType": event.get("detail-type"), |
| 132 | + } |
| 133 | + |
| 134 | + |
| 135 | +def _is_supported_streams(event: dict): |
| 136 | + return event.get("Records", [{}])[0].get("eventSource") in [ |
| 137 | + "aws:kinesis", |
| 138 | + "aws:dynamodb", |
| 139 | + "aws:sqs", |
| 140 | + "aws:s3", |
| 141 | + ] |
| 142 | + |
| 143 | + |
| 144 | +def _parse_streams(event: dict) -> Dict[str, str]: |
| 145 | + """ |
| 146 | + :return: {"triggeredBy": str, "arn": str} |
| 147 | + If has messageId, return also: {"messageId": str} |
| 148 | + """ |
| 149 | + triggered_by = event["Records"][0]["eventSource"].split(":")[1] |
| 150 | + result = {"triggeredBy": triggered_by} |
| 151 | + if triggered_by == "s3": |
| 152 | + result["arn"] = event["Records"][0]["s3"]["bucket"]["arn"] |
| 153 | + result["messageId"] = ( |
| 154 | + event["Records"][0].get("responseElements", {}).get("x-amz-request-id") |
| 155 | + ) |
| 156 | + else: |
| 157 | + result["arn"] = event["Records"][0]["eventSourceARN"] |
| 158 | + if triggered_by == "sqs": |
| 159 | + result.update(_parse_sqs_event(event)) |
| 160 | + elif triggered_by == "kinesis": |
| 161 | + result["messageId"] = safe_get(event, ["Records", 0, "kinesis", "sequenceNumber"]) |
| 162 | + elif triggered_by == "dynamodb": |
| 163 | + result.update(_parse_dynamomdb_event(event)) |
| 164 | + return result |
| 165 | + |
| 166 | + |
| 167 | +def _get_ddb_approx_creation_time_ms(event) -> int: |
| 168 | + return event["Records"][0].get("dynamodb", {}).get("ApproximateCreationDateTime", 0) * 1000 |
| 169 | + |
| 170 | + |
| 171 | +def _parse_dynamomdb_event(event) -> Dict[str, Union[int, List[str]]]: |
| 172 | + creation_time = _get_ddb_approx_creation_time_ms(event) |
| 173 | + mids = [] |
| 174 | + for record in event["Records"]: |
| 175 | + event_name = record.get("eventName") |
| 176 | + if event_name in ("MODIFY", "REMOVE") and record.get("dynamodb", {}).get("Keys"): |
| 177 | + mids.append(md5hash(record["dynamodb"]["Keys"])) |
| 178 | + elif event_name == "INSERT" and record.get("dynamodb", {}).get("NewImage"): |
| 179 | + mids.append(md5hash(record["dynamodb"]["NewImage"])) |
| 180 | + return {MESSAGE_IDS_KEY: mids, TRIGGER_CREATION_TIME_KEY: creation_time} |
| 181 | + |
| 182 | + |
| 183 | +def _parse_sqs_event(event) -> Dict[str, Union[int, List[str]]]: |
| 184 | + mids = [record["messageId"] for record in event["Records"] if record.get("messageId")] |
| 185 | + return {MESSAGE_IDS_KEY: mids} if len(mids) > 1 else {MESSAGE_ID_KEY: mids[0]} |
0 commit comments