Skip to content
Merged
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
2 changes: 1 addition & 1 deletion setup.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ def find_version(*file_paths):
include_package_data=True,
zip_safe=True,
python_requires=">=3.6",
install_requires=["cloudformation-cli>=0.1,<0.2", "docker>=3.7,<5"],
install_requires=["cloudformation-cli>=0.1.10,<0.2", "docker>=3.7,<5"],
Comment thread
johnttompkins marked this conversation as resolved.
entry_points={
"rpdk.v1.languages": [
"python37 = rpdk.python.codegen:Python37LanguagePlugin",
Expand Down
2 changes: 1 addition & 1 deletion src/cloudformation_cli_python_lib/log_delivery.py
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ def setup(
stream_name = f"{request.awsAccountId}-{request.region}"

log_handler = cls._get_existing_logger()
if provider_sess and log_group:
if provider_sess and log_group and request.resourceType:
Comment thread
johnttompkins marked this conversation as resolved.
if log_handler:
# This is a re-used lambda container, log handler is already setup, so
# we just refresh the client with new creds
Expand Down
174 changes: 117 additions & 57 deletions src/cloudformation_cli_python_lib/metrics.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,10 +16,29 @@ def format_dimensions(dimensions: Mapping[str, str]) -> List[Mapping[str, str]]:
return [{"Name": key, "Value": value} for key, value in dimensions.items()]


class MetricPublisher:
def __init__(self, session: SessionProxy, namespace: str) -> None:
self.client = session.client("cloudwatch")
self.namespace = namespace
class MetricsPublisher:
"""A cloudwatch based metric publisher.\
Given a resource type and session, \
this publisher will publish metrics to CloudWatch.\
Can be used with the MetricsPublisherProxy.

Functions:
----------
__init__: Initializes metric publisher with given session and resource type

publish_exception_metric: Publishes an exception based metric

publish_invocation_metric: Publishes a metric related to invocations

publish_duration_metric: Publishes an duration metric

publish_log_delivery_exception_metric: Publishes an log delivery exception metric
"""

def __init__(self, session: SessionProxy, resource_type: str) -> None:
self._client = session.client("cloudwatch")
self._resource_type = resource_type
self._namespace = self._make_namespace(self._resource_type)

def publish_metric( # pylint: disable-msg=too-many-arguments
self,
Expand All @@ -30,8 +49,8 @@ def publish_metric( # pylint: disable-msg=too-many-arguments
timestamp: datetime.datetime,
) -> None:
try:
self.client.put_metric_data(
Namespace=self.namespace,
self._client.put_metric_data(
Namespace=self._namespace,
MetricData=[
{
"MetricName": metric_name.name,
Expand All @@ -46,84 +65,125 @@ def publish_metric( # pylint: disable-msg=too-many-arguments
except ClientError as e:
LOG.error("An error occurred while publishing metrics: %s", str(e))


class MetricsPublisherProxy:
@staticmethod
def _make_namespace(resource_type: str) -> str:
suffix = resource_type.replace("::", "/")
return f"{METRIC_NAMESPACE_ROOT}/{suffix}"

def __init__(self, resource_type: str) -> None:
self.namespace = self._make_namespace(resource_type)
self.resource_type = resource_type
self._publishers: List[MetricPublisher] = []

def add_metrics_publisher(self, session: Optional[SessionProxy]) -> None:
if session:
self._publishers.append(MetricPublisher(session, self.namespace))

def publish_exception_metric(
self, timestamp: datetime.datetime, action: Action, error: Any
) -> None:
dimensions: Mapping[str, str] = {
"DimensionKeyActionType": action.name,
"DimensionKeyExceptionType": str(type(error)),
"DimensionKeyResourceType": self.resource_type,
"DimensionKeyResourceType": self._resource_type,
}
for publisher in self._publishers:
publisher.publish_metric(
metric_name=MetricTypes.HandlerException,
dimensions=dimensions,
unit=StandardUnit.Count,
value=1.0,
timestamp=timestamp,
)
self.publish_metric(
metric_name=MetricTypes.HandlerException,
dimensions=dimensions,
unit=StandardUnit.Count,
value=1.0,
timestamp=timestamp,
)

def publish_invocation_metric(
self, timestamp: datetime.datetime, action: Action
) -> None:
dimensions = {
"DimensionKeyActionType": action.name,
"DimensionKeyResourceType": self.resource_type,
"DimensionKeyResourceType": self._resource_type,
}
for publisher in self._publishers:
publisher.publish_metric(
metric_name=MetricTypes.HandlerInvocationCount,
dimensions=dimensions,
unit=StandardUnit.Count,
value=1.0,
timestamp=timestamp,
)
self.publish_metric(
metric_name=MetricTypes.HandlerInvocationCount,
dimensions=dimensions,
unit=StandardUnit.Count,
value=1.0,
timestamp=timestamp,
)

def publish_duration_metric(
self, timestamp: datetime.datetime, action: Action, milliseconds: float
) -> None:
dimensions = {
"DimensionKeyActionType": action.name,
"DimensionKeyResourceType": self.resource_type,
"DimensionKeyResourceType": self._resource_type,
}
for publisher in self._publishers:
publisher.publish_metric(
metric_name=MetricTypes.HandlerInvocationDuration,
dimensions=dimensions,
unit=StandardUnit.Milliseconds,
value=milliseconds,
timestamp=timestamp,
)

self.publish_metric(
metric_name=MetricTypes.HandlerInvocationDuration,
dimensions=dimensions,
unit=StandardUnit.Milliseconds,
value=milliseconds,
timestamp=timestamp,
)

def publish_log_delivery_exception_metric(
self, timestamp: datetime.datetime, error: Any
) -> None:
dimensions = {
"DimensionKeyActionType": "ProviderLogDelivery",
"DimensionKeyExceptionType": str(type(error)),
"DimensionKeyResourceType": self.resource_type,
"DimensionKeyResourceType": self._resource_type,
}
self.publish_metric(
metric_name=MetricTypes.HandlerException,
dimensions=dimensions,
unit=StandardUnit.Count,
value=1.0,
timestamp=timestamp,
)

@staticmethod
def _make_namespace(resource_type: str) -> str:
suffix = resource_type.replace("::", "/")
return f"{METRIC_NAMESPACE_ROOT}/{suffix}"


class MetricsPublisherProxy:
Comment thread
johnttompkins marked this conversation as resolved.
"""A proxy for publishing metrics to multiple publishers. \
Iterates over available publishers and publishes.

Functions:
----------
add_metrics_publisher: Adds a metrics publisher to the list of publishers

publish_exception_metric: \
Publishes an exception based metric to the list of publishers

publish_invocation_metric: \
Publishes a metric related to invocations to the list of publishers

publish_duration_metric: Publishes a duration metric to the list of publishers

publish_log_delivery_exception_metric: \
Publishes a log delivery exception metric to the list of publishers
"""

def __init__(self) -> None:
self._publishers: List[MetricsPublisher] = []

def add_metrics_publisher(
self, session: Optional[SessionProxy], type_name: Optional[str]
) -> None:
if session and type_name:
publisher = MetricsPublisher(session, type_name)
self._publishers.append(publisher)

def publish_exception_metric(
self, timestamp: datetime.datetime, action: Action, error: Any
) -> None:
for publisher in self._publishers:
publisher.publish_metric(
metric_name=MetricTypes.HandlerException,
dimensions=dimensions,
unit=StandardUnit.Count,
value=1.0,
timestamp=timestamp,
)
publisher.publish_exception_metric(timestamp, action, error)

def publish_invocation_metric(
self, timestamp: datetime.datetime, action: Action
) -> None:
for publisher in self._publishers:
publisher.publish_invocation_metric(timestamp, action)

def publish_duration_metric(
self, timestamp: datetime.datetime, action: Action, milliseconds: float
) -> None:
for publisher in self._publishers:
publisher.publish_duration_metric(timestamp, action, milliseconds)

def publish_log_delivery_exception_metric(
self, timestamp: datetime.datetime, error: Any
) -> None:
for publisher in self._publishers:
publisher.publish_log_delivery_exception_metric(timestamp, error)
16 changes: 6 additions & 10 deletions src/cloudformation_cli_python_lib/resource.py
Original file line number Diff line number Diff line change
Expand Up @@ -148,12 +148,7 @@ def _parse_request(
except Exception as e: # pylint: disable=broad-except
LOG.exception("Invalid request")
raise InvalidRequest(f"{e} ({type(e).__name__})") from e
return (
(caller_sess, provider_sess),
action,
callback_context,
event,
)
return ((caller_sess, provider_sess), action, callback_context, event)

def _cast_resource_request(
self, request: HandlerRequest
Expand Down Expand Up @@ -191,13 +186,14 @@ def print_or_log(message: str) -> None:
try:
sessions, action, callback, event = self._parse_request(event_data)
caller_sess, provider_sess = sessions
ProviderLogHandler.setup(event, provider_sess)
logs_setup = True

request = self._cast_resource_request(event)

metrics = MetricsPublisherProxy(event.resourceType)
metrics.add_metrics_publisher(provider_sess)
metrics = MetricsPublisherProxy()
if event.requestData.providerLogGroupName and provider_sess:
ProviderLogHandler.setup(event, provider_sess)
logs_setup = True
metrics.add_metrics_publisher(provider_sess, event.resourceType)

metrics.publish_invocation_metric(datetime.utcnow(), action)
start_time = datetime.utcnow()
Expand Down
11 changes: 5 additions & 6 deletions src/cloudformation_cli_python_lib/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,9 +47,9 @@ class Credentials:
# pylint: disable=too-many-instance-attributes
@dataclass
class RequestData:
providerLogGroupName: str
logicalResourceId: str
resourceProperties: Mapping[str, Any]
providerLogGroupName: Optional[str] = None
logicalResourceId: Optional[str] = None
systemTags: Optional[Mapping[str, Any]] = None
stackTags: Optional[Mapping[str, Any]] = None
# platform credentials aren't really optional, but this is used to
Expand Down Expand Up @@ -86,13 +86,12 @@ class HandlerRequest:
bearerToken: str
region: str
responseEndpoint: str
resourceType: str
resourceTypeVersion: str
requestData: RequestData
stackId: str
stackId: Optional[str] = None
resourceType: Optional[str] = None
resourceTypeVersion: Optional[str] = None
callbackContext: Optional[MutableMapping[str, Any]] = None
nextToken: Optional[str] = None
requestContext: MutableMapping[str, Any] = field(default_factory=dict)

@classmethod
def deserialize(cls, json_data: MutableMapping[str, Any]) -> "HandlerRequest":
Expand Down
26 changes: 13 additions & 13 deletions tests/lib/metrics_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,15 +7,15 @@
import pytest
from cloudformation_cli_python_lib.interface import Action, MetricTypes, StandardUnit
from cloudformation_cli_python_lib.metrics import (
MetricPublisher,
MetricsPublisher,
MetricsPublisherProxy,
format_dimensions,
)

from botocore.stub import Stubber # pylint: disable=C0411

RESOURCE_TYPE = "Aa::Bb::Cc"
NAMESPACE = MetricsPublisherProxy._make_namespace( # pylint: disable=protected-access
NAMESPACE = MetricsPublisher._make_namespace( # pylint: disable=protected-access
RESOURCE_TYPE
)

Expand Down Expand Up @@ -43,7 +43,7 @@ def test_put_metric_catches_error(mock_session):

mock_session.client.return_value = client

publisher = MetricPublisher(mock_session, NAMESPACE)
publisher = MetricsPublisher(mock_session, NAMESPACE)
dimensions = {
"DimensionKeyActionType": Action.CREATE.name,
"DimensionKeyResourceType": RESOURCE_TYPE,
Expand Down Expand Up @@ -73,8 +73,8 @@ def test_put_metric_catches_error(mock_session):

def test_publish_exception_metric(mock_session):
fake_datetime = datetime(2019, 1, 1)
proxy = MetricsPublisherProxy(RESOURCE_TYPE)
proxy.add_metrics_publisher(mock_session)
proxy = MetricsPublisherProxy()
proxy.add_metrics_publisher(mock_session, RESOURCE_TYPE)
proxy.publish_exception_metric(fake_datetime, Action.CREATE, Exception("fake-err"))
expected_calls = [
call.client("cloudwatch"),
Expand Down Expand Up @@ -103,8 +103,8 @@ def test_publish_exception_metric(mock_session):

def test_publish_invocation_metric(mock_session):
fake_datetime = datetime(2019, 1, 1)
proxy = MetricsPublisherProxy(RESOURCE_TYPE)
proxy.add_metrics_publisher(mock_session)
proxy = MetricsPublisherProxy()
proxy.add_metrics_publisher(mock_session, RESOURCE_TYPE)
proxy.publish_invocation_metric(fake_datetime, Action.CREATE)

expected_calls = [
Expand All @@ -130,8 +130,8 @@ def test_publish_invocation_metric(mock_session):

def test_publish_duration_metric(mock_session):
fake_datetime = datetime(2019, 1, 1)
proxy = MetricsPublisherProxy(RESOURCE_TYPE)
proxy.add_metrics_publisher(mock_session)
proxy = MetricsPublisherProxy()
proxy.add_metrics_publisher(mock_session, RESOURCE_TYPE)
proxy.publish_duration_metric(fake_datetime, Action.CREATE, 100)

expected_calls = [
Expand All @@ -157,8 +157,8 @@ def test_publish_duration_metric(mock_session):

def test_publish_log_delivery_exception_metric(mock_session):
fake_datetime = datetime(2019, 1, 1)
proxy = MetricsPublisherProxy(RESOURCE_TYPE)
proxy.add_metrics_publisher(mock_session)
proxy = MetricsPublisherProxy()
proxy.add_metrics_publisher(mock_session, RESOURCE_TYPE)
proxy.publish_log_delivery_exception_metric(fake_datetime, TypeError("test"))

expected_calls = [
Expand Down Expand Up @@ -190,6 +190,6 @@ def test_publish_log_delivery_exception_metric(mock_session):


def test_metrics_publisher_proxy_add_metrics_publisher_none_safe():
proxy = MetricsPublisherProxy(RESOURCE_TYPE)
proxy.add_metrics_publisher(None)
proxy = MetricsPublisherProxy()
proxy.add_metrics_publisher(None, None)
assert proxy._publishers == [] # pylint: disable=protected-access
4 changes: 2 additions & 2 deletions tests/lib/resource_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@
"responseEndpoint": None,
"resourceType": "AWS::Test::TestModel",
"resourceTypeVersion": "1.0",
"requestContext": {},
"callbackContext": {},
"requestData": {
"callerCredentials": {
"accessKeyId": "IASAYK835GAIFHAHEI23",
Expand Down Expand Up @@ -145,7 +145,7 @@ def test_entrypoint_non_mutating_action():

def test_entrypoint_with_context():
payload = ENTRYPOINT_PAYLOAD.copy()
payload["requestContext"] = {"a": "b"}
payload["callbackContext"] = {"a": "b"}
resource = Resource(TYPE_NAME, Mock())
event = ProgressEvent(
status=OperationStatus.SUCCESS, message="", callbackContext={"c": "d"}
Expand Down
Loading