This is an automated email from the ASF dual-hosted git repository.
absurdfarce pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/cassandra-python-driver.git
The following commit(s) were added to refs/heads/trunk by this push:
new c948fa1e CASSPYTHON-24 Remove Insights support
c948fa1e is described below
commit c948fa1e149989ab9ee164bd2fd025635a94c3e1
Author: absurdfarce <[email protected]>
AuthorDate: Tue Jul 14 21:45:56 2026 -0500
CASSPYTHON-24 Remove Insights support
patch by Bret McGuire; reviewed by Bret McGuire and Brad Schoening
---
cassandra/cluster.py | 81 +-----
cassandra/datastax/insights/__init__.py | 15 -
cassandra/datastax/insights/registry.py | 124 ---------
cassandra/datastax/insights/reporter.py | 219 ---------------
cassandra/datastax/insights/serializers.py | 221 ---------------
cassandra/datastax/insights/util.py | 77 ------
pyproject.toml | 2 +-
.../simulacron/advanced/test_insights.py | 110 --------
tests/integration/standard/test_cluster.py | 5 +-
tests/integration/standard/test_metrics.py | 4 +-
tests/unit/advanced/test_insights.py | 305 ---------------------
11 files changed, 12 insertions(+), 1151 deletions(-)
diff --git a/cassandra/cluster.py b/cassandra/cluster.py
index 9ffefe99..10ce3de1 100644
--- a/cassandra/cluster.py
+++ b/cassandra/cluster.py
@@ -31,10 +31,8 @@ import json
import logging
from warnings import warn
from random import random
-import re
import queue
import socket
-import sys
import time
from threading import Lock, RLock, Thread, Event
import uuid
@@ -83,9 +81,6 @@ from cassandra.marshal import int64_pack
from cassandra.timestamps import MonotonicTimestampGenerator
from cassandra.util import _resolve_contact_points_to_string_map, Version
-from cassandra.datastax.insights.reporter import MonitorReporter
-from cassandra.datastax.insights.util import version_supports_insights
-
from cassandra.datastax.graph import (graph_object_row_factory, GraphOptions,
GraphSON1Serializer,
GraphProtocol, GraphSON2Serializer,
GraphStatement, SimpleGraphStatement,
graph_graphson2_row_factory,
graph_graphson3_row_factory,
@@ -950,34 +945,6 @@ class Cluster(object):
documentation for :meth:`Session.timestamp_generator`.
"""
- monitor_reporting_enabled = True
- """
- A boolean indicating if monitor reporting, which sends gathered data to
- Insights when running against DSE 6.8 and higher.
- """
-
- monitor_reporting_interval = 30
- """
- A boolean indicating if monitor reporting, which sends gathered data to
- Insights when running against DSE 6.8 and higher.
- """
-
- client_id = None
- """
- A UUID that uniquely identifies this Cluster object to Insights. This will
- be generated automatically unless the user provides one.
- """
-
- application_name = ''
- """
- A string identifying this application to Insights.
- """
-
- application_version = ''
- """
- A string identifying this application's version to Insights
- """
-
cloud = None
"""
A dict of the cloud configuration. Example::
@@ -1096,11 +1063,6 @@ class Cluster(object):
no_compact=False,
ssl_context=None,
endpoint_factory=None,
- application_name=None,
- application_version=None,
- monitor_reporting_enabled=True,
- monitor_reporting_interval=30,
- client_id=None,
cloud=None,
column_encryption_policy=None):
"""
@@ -1159,8 +1121,6 @@ class Cluster(object):
raw_contact_points.append(cp if isinstance(cp, tuple) else (cp,
port))
self.endpoints_resolved = [cp for cp in self.contact_points if
isinstance(cp, EndPoint)]
- self._endpoint_map_for_insights = {repr(ep):
'{ip}:{port}'.format(ip=ep.address, port=ep.port)
- for ep in self.endpoints_resolved}
strs_resolved_map =
_resolve_contact_points_to_string_map(raw_contact_points)
self.endpoints_resolved.extend(list(chain(
@@ -1170,14 +1130,14 @@ class Cluster(object):
]
)))
- self._endpoint_map_for_insights.update(
- {key: ['{ip}:{port}'.format(ip=ip, port=port) for ip, port in
value]
- for key, value in strs_resolved_map.items() if value is not None}
- )
-
if contact_points and (not self.endpoints_resolved):
# only want to raise here if the user specified CPs but resolution
failed
- raise UnresolvableContactPoints(self._endpoint_map_for_insights)
+ endpoint_map = {repr(ep): '{ip}:{port}'.format(ip=ep.address,
port=ep.port) for ep in self.endpoints_resolved}
+ endpoint_map.update(
+ {key: ['{ip}:{port}'.format(ip=ip, port=port) for ip, port in
value]
+ for key, value in strs_resolved_map.items() if value is not
None}
+ )
+ raise UnresolvableContactPoints(endpoint_map)
self.compression = compression
@@ -1301,8 +1261,6 @@ class Cluster(object):
self.connect_timeout = connect_timeout
self.prepare_on_all_hosts = prepare_on_all_hosts
self.reprepare_on_up = reprepare_on_up
- self.monitor_reporting_enabled = monitor_reporting_enabled
- self.monitor_reporting_interval = monitor_reporting_interval
self._listeners = set()
self._listener_lock = Lock()
@@ -1352,13 +1310,6 @@ class Cluster(object):
self.status_event_refresh_window,
schema_metadata_enabled, token_metadata_enabled)
- if client_id is None:
- self.client_id = uuid.uuid4()
- if application_name is not None:
- self.application_name = application_name
- if application_version is not None:
- self.application_version = application_version
-
def register_user_type(self, keyspace, user_type, klass):
"""
Registers a class to use to represent a particular user-defined type.
@@ -2473,8 +2424,7 @@ class Session(object):
session_id = None
"""
- A UUID that uniquely identifies this Session to Insights. This will be
- generated automatically.
+ A UUID that uniquely identifies this Session. This will be generated
automatically.
"""
_lock = None
@@ -2527,22 +2477,7 @@ class Session(object):
except AttributeError:
log.info("Unable to set column encryption policy for session")
- if self.cluster.monitor_reporting_enabled:
- cc_host = self.cluster.get_control_connection_host()
- valid_insights_version = (cc_host and
version_supports_insights(cc_host.dse_version))
- if valid_insights_version:
- self._monitor_reporter = MonitorReporter(
- interval_sec=self.cluster.monitor_reporting_interval,
- session=self,
- )
- else:
- if cc_host:
- log.debug('Not starting MonitorReporter thread for
Insights; '
- 'not supported by server version {v} on '
- 'ControlConnection host
{c}'.format(v=cc_host.release_version, c=cc_host))
-
- log.debug('Started Session with client_id {} and session_id
{}'.format(self.cluster.client_id,
-
self.session_id))
+ log.debug('Started Session with session_id {}'.format(self.session_id))
def execute(self, query, parameters=None, timeout=_NOT_SET, trace=False,
custom_payload=None, execution_profile=EXEC_PROFILE_DEFAULT,
diff --git a/cassandra/datastax/insights/__init__.py
b/cassandra/datastax/insights/__init__.py
deleted file mode 100644
index 635f0d9e..00000000
--- a/cassandra/datastax/insights/__init__.py
+++ /dev/null
@@ -1,15 +0,0 @@
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
-#
-# http://www.apache.org/licenses/LICENSE-2.0
-#
-# Unless required by applicable law or agreed to in writing, software
-# distributed under the License is distributed on an "AS IS" BASIS,
-# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-# See the License for the specific language governing permissions and
-# limitations under the License.
diff --git a/cassandra/datastax/insights/registry.py
b/cassandra/datastax/insights/registry.py
deleted file mode 100644
index 523af4dc..00000000
--- a/cassandra/datastax/insights/registry.py
+++ /dev/null
@@ -1,124 +0,0 @@
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
-#
-# http://www.apache.org/licenses/LICENSE-2.0
-#
-# Unless required by applicable law or agreed to in writing, software
-# distributed under the License is distributed on an "AS IS" BASIS,
-# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-# See the License for the specific language governing permissions and
-# limitations under the License.
-
-from collections import OrderedDict
-from warnings import warn
-
-from cassandra.datastax.insights.util import namespace
-
-_NOT_SET = object()
-
-
-def _default_serializer_for_object(obj, policy):
- # the insights server expects an 'options' dict for policy
- # objects, but not for other objects
- if policy:
- return {'type': obj.__class__.__name__,
- 'namespace': namespace(obj.__class__),
- 'options': {}}
- else:
- return {'type': obj.__class__.__name__,
- 'namespace': namespace(obj.__class__)}
-
-
-class InsightsSerializerRegistry(object):
-
- initialized = False
-
- def __init__(self, mapping_dict=None):
- mapping_dict = mapping_dict or {}
- class_order = self._class_topological_sort(mapping_dict)
- self._mapping_dict = OrderedDict(
- ((cls, mapping_dict[cls]) for cls in class_order)
- )
-
- def serialize(self, obj, policy=False, default=_NOT_SET, cls=None):
- try:
- return self._get_serializer(cls if cls is not None else
obj.__class__)(obj)
- except Exception:
- if default is _NOT_SET:
- result = _default_serializer_for_object(obj, policy)
- else:
- result = default
-
- return result
-
- def _get_serializer(self, cls):
- try:
- return self._mapping_dict[cls]
- except KeyError:
- for registered_cls, serializer in self._mapping_dict.items():
- if issubclass(cls, registered_cls):
- return self._mapping_dict[registered_cls]
- raise ValueError
-
- def register(self, cls, serializer):
- self._mapping_dict[cls] = serializer
- self._mapping_dict = OrderedDict(
- ((cls, self._mapping_dict[cls])
- for cls in self._class_topological_sort(self._mapping_dict))
- )
-
- def register_serializer_for(self, cls):
- """
- Parameterized registration helper decorator. Given a class `cls`,
- produces a function that registers the decorated function as a
- serializer for it.
- """
- def decorator(serializer):
- self.register(cls, serializer)
- return serializer
-
- return decorator
-
- @staticmethod
- def _class_topological_sort(classes):
- """
- A simple topological sort for classes. Takes an iterable of class
objects
- and returns a list A of those classes, ordered such that A[X] is never
a
- superclass of A[Y] for X < Y.
-
- This is an inefficient sort, but that's ok because classes are
infrequently
- registered. It's more important that this be maintainable than fast.
-
- We can't use `.sort()` or `sorted()` with a custom `key` -- those
assume
- a total ordering, which we don't have.
- """
- unsorted, sorted_ = list(classes), []
- while unsorted:
- head, tail = unsorted[0], unsorted[1:]
-
- # if head has no subclasses remaining, it can safely go in the list
- if not any(issubclass(x, head) for x in tail):
- sorted_.append(head)
- else:
- # move to the back -- head has to wait until all its subclasses
- # are sorted into the list
- tail.append(head)
-
- unsorted = tail
-
- # check that sort is valid
- for i, head in enumerate(sorted_):
- for after_head_value in sorted_[(i + 1):]:
- if issubclass(after_head_value, head):
- warn('Sorting classes produced an invalid ordering.\n'
- 'In: {classes}\n'
- 'Out: {sorted_}'.format(classes=classes,
sorted_=sorted_))
- return sorted_
-
-
-insights_registry = InsightsSerializerRegistry()
diff --git a/cassandra/datastax/insights/reporter.py
b/cassandra/datastax/insights/reporter.py
deleted file mode 100644
index e3ea5a1c..00000000
--- a/cassandra/datastax/insights/reporter.py
+++ /dev/null
@@ -1,219 +0,0 @@
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
-#
-# http://www.apache.org/licenses/LICENSE-2.0
-#
-# Unless required by applicable law or agreed to in writing, software
-# distributed under the License is distributed on an "AS IS" BASIS,
-# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-# See the License for the specific language governing permissions and
-# limitations under the License.
-
-from collections import Counter
-import datetime
-import json
-import logging
-import multiprocessing
-import random
-import platform
-import socket
-import ssl
-import sys
-from threading import Event, Thread
-import time
-
-from cassandra.policies import HostDistance
-from cassandra.util import ms_timestamp_from_datetime
-from cassandra.datastax.insights.registry import insights_registry
-from cassandra.datastax.insights.serializers import initialize_registry
-
-log = logging.getLogger(__name__)
-
-
-class MonitorReporter(Thread):
-
- def __init__(self, interval_sec, session):
- """
- takes an int indicating interval between requests, a function returning
- the connection to be used, and the timeout per request
- """
- # Thread is an old-style class so we can't super()
- Thread.__init__(self, name='monitor_reporter')
-
- initialize_registry(insights_registry)
-
- self._interval, self._session = interval_sec, session
-
- self._shutdown_event = Event()
- self.daemon = True
- self.start()
-
- def run(self):
- self._send_via_rpc(self._get_startup_data())
-
- # introduce some jitter -- send up to 1/10 of _interval early
- self._shutdown_event.wait(self._interval * random.uniform(.9, 1))
-
- while not self._shutdown_event.is_set():
- start_time = time.time()
-
- self._send_via_rpc(self._get_status_data())
-
- elapsed = time.time() - start_time
- self._shutdown_event.wait(max(self._interval - elapsed, 0.01))
-
- # TODO: redundant with ConnectionHeartbeat.ShutdownException
- class ShutDownException(Exception):
- pass
-
- def _send_via_rpc(self, data):
- try:
- self._session.execute(
- "CALL InsightsRpc.reportInsight(%s)", (json.dumps(data),)
- )
- log.debug('Insights RPC data: {}'.format(data))
- except Exception as e:
- log.debug('Insights RPC send failed with {}'.format(e))
- log.debug('Insights RPC data: {}'.format(data))
-
- def _get_status_data(self):
- cc = self._session.cluster.control_connection
-
- connected_nodes = {
- host.address: {
- 'connections': state['open_count'],
- 'inFlightQueries': state['in_flights']
- }
- for (host, state) in self._session.get_pool_state().items()
- }
-
- return {
- 'metadata': {
- # shared across drivers; never change
- 'name': 'driver.status',
- # format version
- 'insightMappingId': 'v1',
- 'insightType': 'EVENT',
- # since epoch
- 'timestamp':
ms_timestamp_from_datetime(datetime.datetime.utcnow()),
- 'tags': {
- 'language': 'python'
- }
- },
- # // 'clientId', 'sessionId' and 'controlConnection' are mandatory
- # // the rest of the properties are optional
- 'data': {
- # // 'clientId' must be the same as the one provided in the
startup message
- 'clientId': str(self._session.cluster.client_id),
- # // 'sessionId' must be the same as the one provided in the
startup message
- 'sessionId': str(self._session.session_id),
- 'controlConnection': cc._connection.host if cc._connection
else None,
- 'connectedNodes': connected_nodes
- }
- }
-
- def _get_startup_data(self):
- cc = self._session.cluster.control_connection
- try:
- local_ipaddr = cc._connection._socket.getsockname()[0]
- except Exception as e:
- local_ipaddr = None
- log.debug('Unable to get local socket addr from {}:
{}'.format(cc._connection, e))
- hostname = socket.getfqdn()
-
- host_distances_counter = Counter(
- self._session.cluster.profile_manager.distance(host)
- for host in self._session.hosts
- )
- host_distances_dict = {
- 'local': host_distances_counter[HostDistance.LOCAL],
- 'remote': host_distances_counter[HostDistance.REMOTE],
- 'ignored': host_distances_counter[HostDistance.IGNORED]
- }
-
- try:
- compression_type = cc._connection._compression_type
- except AttributeError:
- compression_type = 'NONE'
-
- cert_validation = None
- try:
- if self._session.cluster.ssl_context:
- cert_validation =
self._session.cluster.ssl_context.verify_mode == ssl.CERT_REQUIRED
- elif self._session.cluster.ssl_options:
- cert_validation =
self._session.cluster.ssl_options.get('cert_reqs') == ssl.CERT_REQUIRED
- except Exception as e:
- log.debug('Unable to get the cert validation: {}'.format(e))
-
- uname_info = platform.uname()
-
- return {
- 'metadata': {
- 'name': 'driver.startup',
- 'insightMappingId': 'v1',
- 'insightType': 'EVENT',
- 'timestamp':
ms_timestamp_from_datetime(datetime.datetime.utcnow()),
- 'tags': {
- 'language': 'python'
- },
- },
- 'data': {
- 'driverName': 'DataStax Python Driver',
- 'driverVersion': sys.modules['cassandra'].__version__,
- 'clientId': str(self._session.cluster.client_id),
- 'sessionId': str(self._session.session_id),
- 'applicationName': self._session.cluster.application_name or
'python',
- 'applicationNameWasGenerated': not
self._session.cluster.application_name,
- 'applicationVersion':
self._session.cluster.application_version,
- 'contactPoints':
self._session.cluster._endpoint_map_for_insights,
- 'dataCenters': list(set(h.datacenter for h in
self._session.cluster.metadata.all_hosts()
- if (h.datacenter and
-
self._session.cluster.profile_manager.distance(h) == HostDistance.LOCAL))),
- 'initialControlConnection': cc._connection.host if
cc._connection else None,
- 'protocolVersion': self._session.cluster.protocol_version,
- 'localAddress': local_ipaddr,
- 'hostName': hostname,
- 'executionProfiles':
insights_registry.serialize(self._session.cluster.profile_manager),
- 'configuredConnectionLength': host_distances_dict,
- 'heartbeatInterval':
self._session.cluster.idle_heartbeat_interval,
- 'compression': compression_type.upper() if compression_type
else 'NONE',
- 'reconnectionPolicy':
insights_registry.serialize(self._session.cluster.reconnection_policy),
- 'sslConfigured': {
- 'enabled': bool(self._session.cluster.ssl_options or
self._session.cluster.ssl_context),
- 'certValidation': cert_validation
- },
- 'authProvider': {
- 'type':
(self._session.cluster.auth_provider.__class__.__name__
- if self._session.cluster.auth_provider else
- None)
- },
- 'otherOptions': {
- },
- 'platformInfo': {
- 'os': {
- 'name': uname_info.system,
- 'version': uname_info.release,
- 'arch': uname_info.machine
- },
- 'cpus': {
- 'length': multiprocessing.cpu_count(),
- 'model': platform.processor()
- },
- 'runtime': {
- 'python': sys.version,
- 'event_loop':
self._session.cluster.connection_class.__name__
- }
- },
- 'periodicStatusInterval': self._interval
- }
- }
-
- def stop(self):
- log.debug("Shutting down Monitor Reporter")
- self._shutdown_event.set()
- self.join()
diff --git a/cassandra/datastax/insights/serializers.py
b/cassandra/datastax/insights/serializers.py
deleted file mode 100644
index b1fe0ac5..00000000
--- a/cassandra/datastax/insights/serializers.py
+++ /dev/null
@@ -1,221 +0,0 @@
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
-#
-# http://www.apache.org/licenses/LICENSE-2.0
-#
-# Unless required by applicable law or agreed to in writing, software
-# distributed under the License is distributed on an "AS IS" BASIS,
-# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-# See the License for the specific language governing permissions and
-# limitations under the License.
-
-
-def initialize_registry(insights_registry):
- # This will be called from the cluster module, so we put all this behavior
- # in a function to avoid circular imports
-
- if insights_registry.initialized:
- return False
-
- from cassandra import ConsistencyLevel
- from cassandra.cluster import (
- ExecutionProfile, GraphExecutionProfile,
- ProfileManager, ContinuousPagingOptions,
- EXEC_PROFILE_DEFAULT, EXEC_PROFILE_GRAPH_DEFAULT,
- EXEC_PROFILE_GRAPH_ANALYTICS_DEFAULT,
- EXEC_PROFILE_GRAPH_SYSTEM_DEFAULT,
- _NOT_SET
- )
- from cassandra.datastax.graph import GraphOptions
- from cassandra.datastax.insights.registry import insights_registry
- from cassandra.datastax.insights.util import namespace
- from cassandra.policies import (
- RoundRobinPolicy,
- DCAwareRoundRobinPolicy,
- TokenAwarePolicy,
- WhiteListRoundRobinPolicy,
- HostFilterPolicy,
- ConstantReconnectionPolicy,
- ExponentialReconnectionPolicy,
- RetryPolicy,
- SpeculativeExecutionPolicy,
- ConstantSpeculativeExecutionPolicy,
- WrapperPolicy
- )
-
- import logging
-
- log = logging.getLogger(__name__)
-
- @insights_registry.register_serializer_for(RoundRobinPolicy)
- def round_robin_policy_insights_serializer(policy):
- return {'type': policy.__class__.__name__,
- 'namespace': namespace(policy.__class__),
- 'options': {}}
-
- @insights_registry.register_serializer_for(DCAwareRoundRobinPolicy)
- def dc_aware_round_robin_policy_insights_serializer(policy):
- return {'type': policy.__class__.__name__,
- 'namespace': namespace(policy.__class__),
- 'options': {'local_dc': policy.local_dc,
- 'used_hosts_per_remote_dc':
policy.used_hosts_per_remote_dc}
- }
-
- @insights_registry.register_serializer_for(TokenAwarePolicy)
- def token_aware_policy_insights_serializer(policy):
- return {'type': policy.__class__.__name__,
- 'namespace': namespace(policy.__class__),
- 'options': {'child_policy':
insights_registry.serialize(policy._child_policy,
-
policy=True),
- 'shuffle_replicas': policy.shuffle_replicas}
- }
-
- @insights_registry.register_serializer_for(WhiteListRoundRobinPolicy)
- def whitelist_round_robin_policy_insights_serializer(policy):
- return {'type': policy.__class__.__name__,
- 'namespace': namespace(policy.__class__),
- 'options': {'allowed_hosts': policy._allowed_hosts}
- }
-
- @insights_registry.register_serializer_for(HostFilterPolicy)
- def host_filter_policy_insights_serializer(policy):
- return {
- 'type': policy.__class__.__name__,
- 'namespace': namespace(policy.__class__),
- 'options': {'child_policy':
insights_registry.serialize(policy._child_policy,
-
policy=True),
- 'predicate': policy.predicate.__name__}
- }
-
- @insights_registry.register_serializer_for(ConstantReconnectionPolicy)
- def constant_reconnection_policy_insights_serializer(policy):
- return {'type': policy.__class__.__name__,
- 'namespace': namespace(policy.__class__),
- 'options': {'delay': policy.delay,
- 'max_attempts': policy.max_attempts}
- }
-
- @insights_registry.register_serializer_for(ExponentialReconnectionPolicy)
- def exponential_reconnection_policy_insights_serializer(policy):
- return {'type': policy.__class__.__name__,
- 'namespace': namespace(policy.__class__),
- 'options': {'base_delay': policy.base_delay,
- 'max_delay': policy.max_delay,
- 'max_attempts': policy.max_attempts}
- }
-
- @insights_registry.register_serializer_for(RetryPolicy)
- def retry_policy_insights_serializer(policy):
- return {'type': policy.__class__.__name__,
- 'namespace': namespace(policy.__class__),
- 'options': {}}
-
- @insights_registry.register_serializer_for(SpeculativeExecutionPolicy)
- def speculative_execution_policy_insights_serializer(policy):
- return {'type': policy.__class__.__name__,
- 'namespace': namespace(policy.__class__),
- 'options': {}}
-
-
@insights_registry.register_serializer_for(ConstantSpeculativeExecutionPolicy)
- def constant_speculative_execution_policy_insights_serializer(policy):
- return {'type': policy.__class__.__name__,
- 'namespace': namespace(policy.__class__),
- 'options': {'delay': policy.delay,
- 'max_attempts': policy.max_attempts}
- }
-
- @insights_registry.register_serializer_for(WrapperPolicy)
- def wrapper_policy_insights_serializer(policy):
- return {'type': policy.__class__.__name__,
- 'namespace': namespace(policy.__class__),
- 'options': {
- 'child_policy':
insights_registry.serialize(policy._child_policy,
- policy=True)
- }}
-
- @insights_registry.register_serializer_for(ExecutionProfile)
- def execution_profile_insights_serializer(profile):
- return {
- 'loadBalancing':
insights_registry.serialize(profile.load_balancing_policy,
- policy=True),
- 'retry': insights_registry.serialize(profile.retry_policy,
- policy=True),
- 'readTimeout': profile.request_timeout,
- 'consistency':
ConsistencyLevel.value_to_name.get(profile.consistency_level, None),
- 'serialConsistency':
ConsistencyLevel.value_to_name.get(profile.serial_consistency_level, None),
- 'continuousPagingOptions':
(insights_registry.serialize(profile.continuous_paging_options)
- if (profile.continuous_paging_options
is not None and
- profile.continuous_paging_options
is not _NOT_SET) else
- None),
- 'speculativeExecution':
insights_registry.serialize(profile.speculative_execution_policy),
- 'graphOptions': None
- }
-
- @insights_registry.register_serializer_for(GraphExecutionProfile)
- def graph_execution_profile_insights_serializer(profile):
- rv = insights_registry.serialize(profile, cls=ExecutionProfile)
- rv['graphOptions'] = insights_registry.serialize(profile.graph_options)
- return rv
-
- _EXEC_PROFILE_DEFAULT_KEYS = (EXEC_PROFILE_DEFAULT,
- EXEC_PROFILE_GRAPH_DEFAULT,
- EXEC_PROFILE_GRAPH_SYSTEM_DEFAULT,
- EXEC_PROFILE_GRAPH_ANALYTICS_DEFAULT)
-
- @insights_registry.register_serializer_for(ProfileManager)
- def profile_manager_insights_serializer(manager):
- defaults = {
- # Insights's expected default
- 'default':
insights_registry.serialize(manager.profiles[EXEC_PROFILE_DEFAULT]),
- # remaining named defaults for driver's defaults, including
duplicated default
- 'EXEC_PROFILE_DEFAULT':
insights_registry.serialize(manager.profiles[EXEC_PROFILE_DEFAULT]),
- 'EXEC_PROFILE_GRAPH_DEFAULT':
insights_registry.serialize(manager.profiles[EXEC_PROFILE_GRAPH_DEFAULT]),
- 'EXEC_PROFILE_GRAPH_SYSTEM_DEFAULT': insights_registry.serialize(
- manager.profiles[EXEC_PROFILE_GRAPH_SYSTEM_DEFAULT]
- ),
- 'EXEC_PROFILE_GRAPH_ANALYTICS_DEFAULT':
insights_registry.serialize(
- manager.profiles[EXEC_PROFILE_GRAPH_ANALYTICS_DEFAULT]
- )
- }
- other = {
- key: insights_registry.serialize(value)
- for key, value in manager.profiles.items()
- if key not in _EXEC_PROFILE_DEFAULT_KEYS
- }
- overlapping_keys = set(defaults) & set(other)
- if overlapping_keys:
- log.debug('The following key names overlap default key sentinel
keys '
- 'and these non-default EPs will not be displayed in
Insights '
- ': {}'.format(list(overlapping_keys)))
-
- other.update(defaults)
- return other
-
- @insights_registry.register_serializer_for(GraphOptions)
- def graph_options_insights_serializer(options):
- rv = {
- 'source': options.graph_source,
- 'language': options.graph_language,
- 'graphProtocol': options.graph_protocol
- }
- updates = {k: v.decode('utf-8') for k, v in rv.items()
- if isinstance(v, bytes)}
- rv.update(updates)
- return rv
-
- @insights_registry.register_serializer_for(ContinuousPagingOptions)
- def continuous_paging_options_insights_serializer(paging_options):
- return {
- 'page_unit': paging_options.page_unit,
- 'max_pages': paging_options.max_pages,
- 'max_pages_per_second': paging_options.max_pages_per_second,
- 'max_queue_size': paging_options.max_queue_size
- }
-
- insights_registry.initialized = True
- return True
diff --git a/cassandra/datastax/insights/util.py
b/cassandra/datastax/insights/util.py
deleted file mode 100644
index 0ce96c7e..00000000
--- a/cassandra/datastax/insights/util.py
+++ /dev/null
@@ -1,77 +0,0 @@
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
-#
-# http://www.apache.org/licenses/LICENSE-2.0
-#
-# Unless required by applicable law or agreed to in writing, software
-# distributed under the License is distributed on an "AS IS" BASIS,
-# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-# See the License for the specific language governing permissions and
-# limitations under the License.
-
-import logging
-import traceback
-from warnings import warn
-
-from cassandra.util import Version
-
-
-DSE_60 = Version('6.0.0')
-DSE_51_MIN_SUPPORTED = Version('5.1.13')
-DSE_60_MIN_SUPPORTED = Version('6.0.5')
-
-
-log = logging.getLogger(__name__)
-
-
-def namespace(cls):
- """
- Best-effort method for getting the namespace in which a class is defined.
- """
- try:
- # __module__ can be None
- module = cls.__module__ or ''
- except Exception:
- warn("Unable to obtain namespace for {cls} for Insights, returning ''.
"
- "Exception: \n{e}".format(e=traceback.format_exc(), cls=cls))
- module = ''
-
- module_internal_namespace = _module_internal_namespace_or_emtpy_string(cls)
- if module_internal_namespace:
- return '.'.join((module, module_internal_namespace))
- return module
-
-
-def _module_internal_namespace_or_emtpy_string(cls):
- """
- Best-effort method for getting the module-internal namespace in which a
- class is defined -- i.e. the namespace _inside_ the module.
- """
- try:
- qualname = cls.__qualname__
- except AttributeError:
- return ''
-
- return '.'.join(
- # the last segment is the name of the class -- use everything else
- qualname.split('.')[:-1]
- )
-
-
-def version_supports_insights(dse_version):
- if dse_version:
- try:
- dse_version = Version(dse_version)
- return (DSE_51_MIN_SUPPORTED <= dse_version < DSE_60
- or
- DSE_60_MIN_SUPPORTED <= dse_version)
- except Exception:
- warn("Unable to check version {v} for Insights compatibility,
returning False. "
- "Exception: \n{e}".format(e=traceback.format_exc(),
v=dse_version))
-
- return False
diff --git a/pyproject.toml b/pyproject.toml
index 0af1e770..9f237f0e 100644
--- a/pyproject.toml
+++ b/pyproject.toml
@@ -43,7 +43,7 @@ changelog =
"https://github.com/apache/cassandra-python-driver/blob/trunk/CHANGE
[tool.setuptools.packages.find]
include = ['cassandra', 'cassandra.io', 'cassandra.cqlengine',
'cassandra.graph',
-'cassandra.datastax', 'cassandra.datastax.insights',
'cassandra.datastax.graph',
+'cassandra.datastax', 'cassandra.datastax.graph',
'cassandra.datastax.graph.fluent', 'cassandra.datastax.cloud',
"cassandra.column_encryption"]
diff --git a/tests/integration/simulacron/advanced/test_insights.py
b/tests/integration/simulacron/advanced/test_insights.py
deleted file mode 100644
index 07005a47..00000000
--- a/tests/integration/simulacron/advanced/test_insights.py
+++ /dev/null
@@ -1,110 +0,0 @@
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
-#
-# http://www.apache.org/licenses/LICENSE-2.0
-#
-# Unless required by applicable law or agreed to in writing, software
-# distributed under the License is distributed on an "AS IS" BASIS,
-# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-# See the License for the specific language governing permissions and
-# limitations under the License.
-import unittest
-
-import time
-import json
-import re
-
-from cassandra.cluster import Cluster
-from cassandra.datastax.insights.util import version_supports_insights
-
-from tests.integration import requiressimulacron, requiredse, DSE_VERSION
-from tests.integration.simulacron import DseSimulacronCluster, PROTOCOL_VERSION
-from tests.integration.simulacron.utils import SimulacronClient, GetLogsQuery,
ClearLogsQuery
-
-
-@requiredse
-@requiressimulacron
[email protected](DSE_VERSION and
version_supports_insights(str(DSE_VERSION)), 'DSE {} does not support
insights'.format(DSE_VERSION))
-class InsightsTests(DseSimulacronCluster):
- """
- Tests insights integration
-
- @since 3.18
- @jira_ticket PYTHON-1047
- @expected_result startup and status messages are sent
- """
-
- connect = False
-
- def tearDown(self):
- if self.cluster:
- self.cluster.shutdown()
-
- @staticmethod
- def _get_node_logs(raw_data):
- return list(filter(lambda q: q['type'] == 'QUERY' and
q['query'].startswith('CALL InsightsRpc.reportInsight'),
-
json.loads(raw_data)['data_centers'][0]['nodes'][0]['queries']))
-
- @staticmethod
- def _parse_data(data, index=0):
- return json.loads(re.match(
- r"CALL InsightsRpc.reportInsight\('(.+)'\)",
- data[index]['frame']['message']['query']).group(1))
-
- def test_startup_message(self):
- self.cluster = Cluster(protocol_version=PROTOCOL_VERSION,
compression=False)
- self.session = self.cluster.connect(wait_for_all_pools=True)
-
- time.sleep(1) # wait the monitor thread is started
- response = SimulacronClient().submit_request(GetLogsQuery())
- self.assertTrue('CALL InsightsRpc.reportInsight' in response)
-
- node_queries = self._get_node_logs(response)
- self.assertEqual(1, len(node_queries))
- self.assertTrue(node_queries, "RPC query not found")
-
- message = self._parse_data(node_queries)
-
- self.assertEqual(message['metadata']['name'], 'driver.startup')
- self.assertEqual(message['data']['initialControlConnection'],
- self.cluster.control_connection._connection.host)
- self.assertEqual(message['data']['sessionId'],
str(self.session.session_id))
- self.assertEqual(message['data']['clientId'],
str(self.cluster.client_id))
- self.assertEqual(message['data']['compression'], 'NONE')
-
- def test_status_message(self):
- SimulacronClient().submit_request(ClearLogsQuery())
-
- self.cluster = Cluster(protocol_version=PROTOCOL_VERSION,
compression=False, monitor_reporting_interval=1)
- self.session = self.cluster.connect(wait_for_all_pools=True)
-
- time.sleep(1.1)
- response = SimulacronClient().submit_request(GetLogsQuery())
- self.assertTrue('CALL InsightsRpc.reportInsight' in response)
-
- node_queries = self._get_node_logs(response)
- self.assertEqual(2, len(node_queries))
- self.assertTrue(node_queries, "RPC query not found")
-
- message = self._parse_data(node_queries, 1)
-
- self.assertEqual(message['metadata']['name'], 'driver.status')
- self.assertEqual(message['data']['controlConnection'],
- self.cluster.control_connection._connection.host)
- self.assertEqual(message['data']['sessionId'],
str(self.session.session_id))
- self.assertEqual(message['data']['clientId'],
str(self.cluster.client_id))
- self.assertEqual(message['metadata']['insightType'], 'EVENT')
-
- def test_monitor_disabled(self):
- SimulacronClient().submit_request(ClearLogsQuery())
-
- self.cluster = Cluster(protocol_version=PROTOCOL_VERSION,
compression=False, monitor_reporting_enabled=False)
- self.session = self.cluster.connect(wait_for_all_pools=True)
-
- response = SimulacronClient().submit_request(GetLogsQuery())
- self.assertFalse('CALL InsightsRpc.reportInsight' in response)
diff --git a/tests/integration/standard/test_cluster.py
b/tests/integration/standard/test_cluster.py
index c6fc2a71..3dd4bda1 100644
--- a/tests/integration/standard/test_cluster.py
+++ b/tests/integration/standard/test_cluster.py
@@ -754,8 +754,7 @@ class ClusterTests(unittest.TestCase):
def test_idle_heartbeat(self):
interval = 2
- cluster = TestCluster(idle_heartbeat_interval=interval,
- monitor_reporting_enabled=False)
+ cluster = TestCluster(idle_heartbeat_interval=interval)
if PROTOCOL_VERSION < 3:
cluster.set_core_connections_per_host(HostDistance.LOCAL, 1)
session = cluster.connect(wait_for_all_pools=True)
@@ -877,7 +876,7 @@ class ClusterTests(unittest.TestCase):
RoundRobinPolicy(), lambda host: host.address == CASSANDRA_IP
)
)
- with TestCluster(execution_profiles={'node1': node1},
monitor_reporting_enabled=False) as cluster:
+ with TestCluster(execution_profiles={'node1': node1}) as cluster:
session = cluster.connect(wait_for_all_pools=True)
# default is DCA RR for all hosts
diff --git a/tests/integration/standard/test_metrics.py
b/tests/integration/standard/test_metrics.py
index c33ea265..44cd712d 100644
--- a/tests/integration/standard/test_metrics.py
+++ b/tests/integration/standard/test_metrics.py
@@ -258,13 +258,11 @@ class
MetricsNamespaceTest(BasicSharedKeyspaceUnitTestCaseRF3WM):
"""
cluster2 = TestCluster(
metrics_enabled=True,
- monitor_reporting_enabled=False,
execution_profiles={EXEC_PROFILE_DEFAULT:
ExecutionProfile(retry_policy=FallthroughRetryPolicy())}
)
cluster3 = TestCluster(
metrics_enabled=True,
- monitor_reporting_enabled=False,
execution_profiles={EXEC_PROFILE_DEFAULT:
ExecutionProfile(retry_policy=FallthroughRetryPolicy())}
)
@@ -351,7 +349,7 @@ class MetricsRequestSize(BasicExistingKeyspaceUnitTestCase):
@classmethod
def setUpClass(cls):
- cls.common_setup(1, keyspace_creation=False,
monitor_reporting_enabled=False)
+ cls.common_setup(1, keyspace_creation=False)
def wait_for_count(self, ra, expected_count, error=False):
for _ in range(10):
diff --git a/tests/unit/advanced/test_insights.py
b/tests/unit/advanced/test_insights.py
deleted file mode 100644
index aae17879..00000000
--- a/tests/unit/advanced/test_insights.py
+++ /dev/null
@@ -1,305 +0,0 @@
-# Licensed to the Apache Software Foundation (ASF) under one
-# or more contributor license agreements. See the NOTICE file
-# distributed with this work for additional information
-# regarding copyright ownership. The ASF licenses this file
-# to you under the Apache License, Version 2.0 (the
-# "License"); you may not use this file except in compliance
-# with the License. You may obtain a copy of the License at
-#
-# http://www.apache.org/licenses/LICENSE-2.0
-#
-# Unless required by applicable law or agreed to in writing, software
-# distributed under the License is distributed on an "AS IS" BASIS,
-# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-# See the License for the specific language governing permissions and
-# limitations under the License.
-
-
-import unittest
-
-import logging
-from unittest.mock import sentinel
-
-from cassandra import ConsistencyLevel
-from cassandra.cluster import (
- ExecutionProfile, GraphExecutionProfile, ProfileManager,
- GraphAnalyticsExecutionProfile,
- EXEC_PROFILE_DEFAULT, EXEC_PROFILE_GRAPH_DEFAULT,
- EXEC_PROFILE_GRAPH_ANALYTICS_DEFAULT,
- EXEC_PROFILE_GRAPH_SYSTEM_DEFAULT
-)
-from cassandra.datastax.graph.query import GraphOptions
-from cassandra.datastax.insights.registry import insights_registry
-from cassandra.datastax.insights.serializers import initialize_registry
-from cassandra.policies import (
- RoundRobinPolicy,
- LoadBalancingPolicy,
- DCAwareRoundRobinPolicy,
- TokenAwarePolicy,
- WhiteListRoundRobinPolicy,
- HostFilterPolicy,
- ConstantReconnectionPolicy,
- ExponentialReconnectionPolicy,
- RetryPolicy,
- SpeculativeExecutionPolicy,
- ConstantSpeculativeExecutionPolicy,
- WrapperPolicy
-)
-
-
-log = logging.getLogger(__name__)
-
-initialize_registry(insights_registry)
-
-
-class TestGetConfig(unittest.TestCase):
-
- def test_invalid_object(self):
- class NoConfAsDict(object):
- pass
-
- obj = NoConfAsDict()
-
- ns = 'tests.unit.advanced.test_insights'
- ns += '.TestGetConfig.test_invalid_object.<locals>'
-
- # no default
- # ... as a policy
- self.assertEqual(insights_registry.serialize(obj, policy=True),
- {'type': 'NoConfAsDict',
- 'namespace': ns,
- 'options': {}})
- # ... not as a policy (default)
- self.assertEqual(insights_registry.serialize(obj),
- {'type': 'NoConfAsDict',
- 'namespace': ns,
- })
- # with default
- self.assertIs(insights_registry.serialize(obj,
default=sentinel.attr_err_default),
- sentinel.attr_err_default)
-
- def test_successful_return(self):
-
- class SuperclassSentinel(object):
- pass
-
- class SubclassSentinel(SuperclassSentinel):
- pass
-
- @insights_registry.register_serializer_for(SuperclassSentinel)
- def superclass_sentinel_serializer(obj):
- return sentinel.serialized_superclass
-
- self.assertIs(insights_registry.serialize(SuperclassSentinel()),
- sentinel.serialized_superclass)
- self.assertIs(insights_registry.serialize(SubclassSentinel()),
- sentinel.serialized_superclass)
-
- # with default -- same behavior
- self.assertIs(insights_registry.serialize(SubclassSentinel(),
default=object()),
- sentinel.serialized_superclass)
-
-
-class TestConfigAsDict(unittest.TestCase):
-
- # graph/query.py
- def test_graph_options(self):
- self.maxDiff = None
-
- go = GraphOptions(graph_name='name_for_test',
- graph_source='source_for_test',
- graph_language='lang_for_test',
- graph_protocol='protocol_for_test',
- graph_read_consistency_level=ConsistencyLevel.ANY,
- graph_write_consistency_level=ConsistencyLevel.ONE,
- graph_invalid_option='invalid')
-
- log.debug(go._graph_options)
-
- self.assertEqual(
- insights_registry.serialize(go),
- {'source': 'source_for_test',
- 'language': 'lang_for_test',
- 'graphProtocol': 'protocol_for_test',
- # no graph_invalid_option
- }
- )
-
- # cluster.py
- def test_execution_profile(self):
- self.maxDiff = None
- self.assertEqual(
- insights_registry.serialize(ExecutionProfile()),
- {'consistency': 'LOCAL_ONE',
- 'continuousPagingOptions': None,
- 'loadBalancing': {'namespace': 'cassandra.policies',
- 'options': {'child_policy': {'namespace':
'cassandra.policies',
- 'options':
{'local_dc': '',
-
'used_hosts_per_remote_dc': 0},
- 'type':
'DCAwareRoundRobinPolicy'},
- 'shuffle_replicas': False},
- 'type': 'TokenAwarePolicy'},
- 'readTimeout': 10.0,
- 'retry': {'namespace': 'cassandra.policies', 'options': {},
'type': 'RetryPolicy'},
- 'serialConsistency': None,
- 'speculativeExecution': {'namespace': 'cassandra.policies',
- 'options': {}, 'type':
'NoSpeculativeExecutionPolicy'},
- 'graphOptions': None
- }
- )
-
- def test_graph_execution_profile(self):
- self.maxDiff = None
- self.assertEqual(
- insights_registry.serialize(GraphExecutionProfile()),
- {'consistency': 'LOCAL_ONE',
- 'continuousPagingOptions': None,
- 'loadBalancing': {'namespace': 'cassandra.policies',
- 'options': {'child_policy': {'namespace':
'cassandra.policies',
- 'options':
{'local_dc': '',
-
'used_hosts_per_remote_dc': 0},
- 'type':
'DCAwareRoundRobinPolicy'},
- 'shuffle_replicas': False},
- 'type': 'TokenAwarePolicy'},
- 'readTimeout': 30.0,
- 'retry': {'namespace': 'cassandra.policies', 'options': {},
'type': 'NeverRetryPolicy'},
- 'serialConsistency': None,
- 'speculativeExecution': {'namespace': 'cassandra.policies',
- 'options': {}, 'type':
'NoSpeculativeExecutionPolicy'},
- 'graphOptions': {'graphProtocol': None,
- 'language': 'gremlin-groovy',
- 'source': 'g'},
- }
- )
-
- def test_graph_analytics_execution_profile(self):
- self.maxDiff = None
- self.assertEqual(
- insights_registry.serialize(GraphAnalyticsExecutionProfile()),
- {'consistency': 'LOCAL_ONE',
- 'continuousPagingOptions': None,
- 'loadBalancing': {'namespace': 'cassandra.policies',
- 'options': {'child_policy': {'namespace':
'cassandra.policies',
- 'options':
{'child_policy': {'namespace': 'cassandra.policies',
-
'options': {'local_dc': '',
-
'used_hosts_per_remote_dc': 0},
-
'type': 'DCAwareRoundRobinPolicy'},
-
'shuffle_replicas': False},
- 'type':
'TokenAwarePolicy'}},
- 'type': 'DefaultLoadBalancingPolicy'},
- 'readTimeout': 604800.0,
- 'retry': {'namespace': 'cassandra.policies', 'options': {},
'type': 'NeverRetryPolicy'},
- 'serialConsistency': None,
- 'speculativeExecution': {'namespace': 'cassandra.policies',
- 'options': {}, 'type':
'NoSpeculativeExecutionPolicy'},
- 'graphOptions': {'graphProtocol': None,
- 'language': 'gremlin-groovy',
- 'source': 'a'},
- }
- )
-
- # policies.py
- def test_DC_aware_round_robin_policy(self):
- self.assertEqual(
- insights_registry.serialize(DCAwareRoundRobinPolicy()),
- {'namespace': 'cassandra.policies',
- 'options': {'local_dc': '', 'used_hosts_per_remote_dc': 0},
- 'type': 'DCAwareRoundRobinPolicy'}
- )
- self.assertEqual(
-
insights_registry.serialize(DCAwareRoundRobinPolicy(local_dc='fake_local_dc',
-
used_hosts_per_remote_dc=15)),
- {'namespace': 'cassandra.policies',
- 'options': {'local_dc': 'fake_local_dc',
'used_hosts_per_remote_dc': 15},
- 'type': 'DCAwareRoundRobinPolicy'}
- )
-
- def test_token_aware_policy(self):
- self.assertEqual(
-
insights_registry.serialize(TokenAwarePolicy(child_policy=LoadBalancingPolicy())),
- {'namespace': 'cassandra.policies',
- 'options': {'child_policy': {'namespace': 'cassandra.policies',
- 'options': {},
- 'type': 'LoadBalancingPolicy'},
- 'shuffle_replicas': False},
- 'type': 'TokenAwarePolicy'}
- )
-
- def test_whitelist_round_robin_policy(self):
- self.assertEqual(
-
insights_registry.serialize(WhiteListRoundRobinPolicy(['127.0.0.3'])),
- {'namespace': 'cassandra.policies',
- 'options': {'allowed_hosts': ('127.0.0.3',)},
- 'type': 'WhiteListRoundRobinPolicy'}
- )
-
- def test_host_filter_policy(self):
- def my_predicate(s):
- return False
-
- self.assertEqual(
-
insights_registry.serialize(HostFilterPolicy(LoadBalancingPolicy(),
my_predicate)),
- {'namespace': 'cassandra.policies',
- 'options': {'child_policy': {'namespace': 'cassandra.policies',
- 'options': {},
- 'type': 'LoadBalancingPolicy'},
- 'predicate': 'my_predicate'},
- 'type': 'HostFilterPolicy'}
- )
-
- def test_constant_reconnection_policy(self):
- self.assertEqual(
- insights_registry.serialize(ConstantReconnectionPolicy(3, 200)),
- {'type': 'ConstantReconnectionPolicy',
- 'namespace': 'cassandra.policies',
- 'options': {'delay': 3, 'max_attempts': 200}
- }
- )
-
- def test_exponential_reconnection_policy(self):
- self.assertEqual(
- insights_registry.serialize(ExponentialReconnectionPolicy(4, 100,
10)),
- {'type': 'ExponentialReconnectionPolicy',
- 'namespace': 'cassandra.policies',
- 'options': {'base_delay': 4, 'max_delay': 100, 'max_attempts':
10}
- }
- )
-
- def test_retry_policy(self):
- self.assertEqual(
- insights_registry.serialize(RetryPolicy()),
- {'type': 'RetryPolicy',
- 'namespace': 'cassandra.policies',
- 'options': {}
- }
- )
-
- def test_spec_exec_policy(self):
- self.assertEqual(
- insights_registry.serialize(SpeculativeExecutionPolicy()),
- {'type': 'SpeculativeExecutionPolicy',
- 'namespace': 'cassandra.policies',
- 'options': {}
- }
- )
-
- def test_constant_spec_exec_policy(self):
- self.assertEqual(
-
insights_registry.serialize(ConstantSpeculativeExecutionPolicy(100, 101)),
- {'type': 'ConstantSpeculativeExecutionPolicy',
- 'namespace': 'cassandra.policies',
- 'options': {'delay': 100,
- 'max_attempts': 101}
- }
- )
-
- def test_wrapper_policy(self):
- self.assertEqual(
- insights_registry.serialize(WrapperPolicy(LoadBalancingPolicy())),
- {'namespace': 'cassandra.policies',
- 'options': {'child_policy': {'namespace': 'cassandra.policies',
- 'options': {},
- 'type': 'LoadBalancingPolicy'}
- },
- 'type': 'WrapperPolicy'}
- )
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]