diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index de21a2fb9..753e3bc28 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -601,7 +601,6 @@ jobs: env: CSP_TEST_SKIP_EXAMPLES: "1" - #################################################### #..................................................# #..|########|..|########|..../####\....|########|..# @@ -713,7 +712,75 @@ jobs: ########################################################################################################### # Test Service Adapters # ########################################################################################################### - # Coming soon! + test_adapters: + needs: + - initialize + - build + + permissions: + contents: read + + timeout-minutes: 30 + + strategy: + matrix: + os: + - ubuntu-24.04 + python-version: + - 3.11 + adapter: + - kafka + + runs-on: ${{ matrix.os }} + + steps: + - name: Checkout + uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + with: + submodules: recursive + persist-credentials: false + + - name: Set up Python ${{ matrix.python-version }} + uses: ./.github/actions/setup-python + with: + version: '${{ matrix.python-version }}' + cibuildwheel: false + + - name: Install python dependencies + run: make requirements + + - name: Install test dependencies + shell: bash + run: sudo apt-get update && sudo apt-get install -y graphviz + + # Download artifact + - name: Download wheel + uses: actions/download-artifact@3e5f45b2cfb9172054b4087a40e8e0b5a5461e7c # v8.0.1 + with: + name: csp-dist-${{ runner.os }}-${{ runner.arch }}-${{ matrix.python-version }} + + - name: Install wheel + run: | + python -m pip install -U *manylinux*.whl + python -m pip install -U --no-deps *manylinux*.whl --target . + + # Services declare healthchecks, so --wait blocks until the broker actually accepts connections + - name: Spin up adapter service + run: make dockerup ADAPTER=${{ matrix.adapter }} DOCKERARGS="--wait --wait-timeout 180" + + # Run tests + - name: Setup test flags + shell: bash + env: + ADAPTER: ${{ matrix.adapter }} + run: echo "CSP_TEST_${ADAPTER^^}=1" >> "$GITHUB_ENV" + + - name: Python Test Steps + run: make test-py TEST_ARGS="-k ${{ matrix.adapter }}" + + - name: Spin down adapter service + run: make dockerdown ADAPTER=${{ matrix.adapter }} + if: ${{ always() }} ############################################################################################ #..........................................................................................# @@ -728,7 +795,6 @@ jobs: ############################################################################################ # Upload Release Artifacts # ############################################################################################ - # only publish artifacts on tags, but otherwise this always runs # Note this whole workflow only triggers on release tags (e.g. "v0.1.0") publish_release_artifacts: @@ -740,6 +806,7 @@ jobs: - test - test_sdist - test_dependencies + - test_adapters if: startsWith(github.ref, 'refs/tags/v') runs-on: ubuntu-24.04 diff --git a/Makefile b/Makefile index b55b52d2e..4f5d49a5a 100644 --- a/Makefile +++ b/Makefile @@ -117,7 +117,9 @@ tests: test .PHONY: dockerup dockerps dockerdown initpodmanmac ADAPTER := kafka -DOCKER := podman +# Prefer docker, fall back to podman-compose; override with DOCKER_COMPOSE=... +DOCKER_COMPOSE := $(shell command -v docker >/dev/null 2>&1 && echo "docker compose" || echo "podman-compose") +DOCKERARGS := initpodmanmac: podman machine stop @@ -125,13 +127,13 @@ initpodmanmac: podman machine start dockerup: ## spin up docker compose services for adapter testing - $(DOCKER) compose -f ci/$(ADAPTER)/docker-compose.yml up -d + $(DOCKER_COMPOSE) -f ci/$(ADAPTER)/docker-compose.yml up -d $(DOCKERARGS) -dockerps: ## spin up docker compose services for adapter testing - $(DOCKER) compose -f ci/$(ADAPTER)/docker-compose.yml ps +dockerps: ## get status of current docker compose services + $(DOCKER_COMPOSE) -f ci/$(ADAPTER)/docker-compose.yml ps -dockerdown: ## spin up docker compose services for adapter testing - $(DOCKER) compose -f ci/$(ADAPTER)/docker-compose.yml down +dockerdown: ## spin down docker compose services for adapter testing + $(DOCKER_COMPOSE) -f ci/$(ADAPTER)/docker-compose.yml down ########### # VERSION # @@ -222,9 +224,11 @@ dependencies-fedora: ## install dependencies for linux - note that zip is neede dependencies-vcpkg: ## install dependencies via vcpkg cd vcpkg && ./bootstrap-vcpkg.sh && ./vcpkg install +# curl.exe ships with Windows; installing it from choco breaks whenever a new version is +# approved on the feed before the package is actually downloadable dependencies-win: ## install dependencies via windows choco install cmake --version=3.31.6 --allow-downgrade - choco install curl winflexbison ninja unzip --no-progress -y + choco install winflexbison ninja unzip --no-progress -y ############################################################################################ # Thanks to Francoise at marmelab.com for this diff --git a/ci/kafka/docker-compose.yml b/ci/kafka/docker-compose.yml index d2674945f..a15cbcce0 100644 --- a/ci/kafka/docker-compose.yml +++ b/ci/kafka/docker-compose.yml @@ -1,181 +1,48 @@ -# https://docs.confluent.io/platform/current/platform-quickstart.html -# https://raw.githubusercontent.com/confluentinc/cp-all-in-one/7.5.3-post/cp-all-in-one-kraft/docker-compose.yml +# https://github.com/conduktor/kafka-stack-docker-compose --- -version: '2' services: - zookeeper: + zoo1: image: confluentinc/cp-zookeeper:7.5.3 - hostname: zookeeper - container_name: zookeeper + hostname: zoo1 + container_name: zoo1 ports: - - "2181:2181" + - "127.0.0.1:2181:2181" environment: ZOOKEEPER_CLIENT_PORT: 2181 - ZOOKEEPER_TICK_TIME: 2000 - - broker: - image: confluentinc/cp-server:7.5.3 - hostname: broker - container_name: broker - depends_on: - - zookeeper + ZOOKEEPER_SERVER_ID: 1 + ZOOKEEPER_SERVERS: zoo1:2888:3888 + healthcheck: + # cub is shipped in the image; the 4lw "ruok" probe is not whitelisted by default in ZK 3.5+ + test: ["CMD", "cub", "zk-ready", "localhost:2181", "5"] + interval: 5s + timeout: 10s + retries: 24 + + kafka1: + image: confluentinc/cp-kafka:7.5.3 + hostname: kafka1 + container_name: kafka1 ports: - - "9092:9092" - - "9101:9101" + - "127.0.0.1:9092:9092" environment: + KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka1:19092,EXTERNAL://${DOCKER_HOST_IP:-127.0.0.1}:9092 + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT + KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL + KAFKA_ZOOKEEPER_CONNECT: "zoo1:2181" KAFKA_BROKER_ID: 1 - KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' - KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT - KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://broker:29092,PLAINTEXT_HOST://localhost:9092 - KAFKA_METRIC_REPORTERS: io.confluent.metrics.reporter.ConfluentMetricsReporter + KAFKA_LOG4J_LOGGERS: "kafka.controller=INFO,kafka.producer.async.DefaultEventHandler=INFO,state.change.logger=INFO" KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 - KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 - KAFKA_CONFLUENT_LICENSE_TOPIC_REPLICATION_FACTOR: 1 - KAFKA_CONFLUENT_BALANCER_TOPIC_REPLICATION_FACTOR: 1 - KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 - KAFKA_JMX_PORT: 9101 - KAFKA_JMX_HOSTNAME: localhost - KAFKA_CONFLUENT_SCHEMA_REGISTRY_URL: http://schema-registry:8081 - CONFLUENT_METRICS_REPORTER_BOOTSTRAP_SERVERS: broker:29092 - CONFLUENT_METRICS_REPORTER_TOPIC_REPLICAS: 1 - CONFLUENT_METRICS_ENABLE: 'true' - CONFLUENT_SUPPORT_CUSTOMER_ID: 'anonymous' - KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'true' - TOPIC_AUTO_CREATE: 'true' - - schema-registry: - image: confluentinc/cp-schema-registry:7.5.3 - hostname: schema-registry - container_name: schema-registry - depends_on: - - broker - ports: - - "8081:8081" - environment: - SCHEMA_REGISTRY_HOST_NAME: schema-registry - SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: 'broker:29092' - SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:8081 - - connect: - image: cnfldemos/cp-server-connect-datagen:0.6.2-7.5.0 - hostname: connect - container_name: connect - depends_on: - - broker - - schema-registry - ports: - - "8083:8083" - environment: - CONNECT_BOOTSTRAP_SERVERS: 'broker:29092' - CONNECT_REST_ADVERTISED_HOST_NAME: connect - CONNECT_GROUP_ID: compose-connect-group - CONNECT_CONFIG_STORAGE_TOPIC: docker-connect-configs - CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR: 1 - CONNECT_OFFSET_FLUSH_INTERVAL_MS: 10000 - CONNECT_OFFSET_STORAGE_TOPIC: docker-connect-offsets - CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: 1 - CONNECT_STATUS_STORAGE_TOPIC: docker-connect-status - CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: 1 - CONNECT_KEY_CONVERTER: org.apache.kafka.connect.storage.StringConverter - CONNECT_VALUE_CONVERTER: io.confluent.connect.avro.AvroConverter - CONNECT_VALUE_CONVERTER_SCHEMA_REGISTRY_URL: http://schema-registry:8081 - # CLASSPATH required due to CC-2422 - CLASSPATH: /usr/share/java/monitoring-interceptors/monitoring-interceptors-7.5.3.jar - CONNECT_PRODUCER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringProducerInterceptor" - CONNECT_CONSUMER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringConsumerInterceptor" - CONNECT_PLUGIN_PATH: "/usr/share/java,/usr/share/confluent-hub-components" - CONNECT_LOG4J_LOGGERS: org.apache.zookeeper=ERROR,org.I0Itec.zkclient=ERROR,org.reflections=ERROR - - control-center: - image: confluentinc/cp-enterprise-control-center:7.5.3 - hostname: control-center - container_name: control-center - depends_on: - - broker - - schema-registry - - connect - - ksqldb-server - ports: - - "9021:9021" - environment: - CONTROL_CENTER_BOOTSTRAP_SERVERS: 'broker:29092' - CONTROL_CENTER_CONNECT_CONNECT-DEFAULT_CLUSTER: 'connect:8083' - CONTROL_CENTER_KSQL_KSQLDB1_URL: "http://ksqldb-server:8088" - CONTROL_CENTER_KSQL_KSQLDB1_ADVERTISED_URL: "http://localhost:8088" - CONTROL_CENTER_SCHEMA_REGISTRY_URL: "http://schema-registry:8081" - CONTROL_CENTER_REPLICATION_FACTOR: 1 - CONTROL_CENTER_INTERNAL_TOPICS_PARTITIONS: 1 - CONTROL_CENTER_MONITORING_INTERCEPTOR_TOPIC_PARTITIONS: 1 - CONFLUENT_METRICS_TOPIC_REPLICATION: 1 - PORT: 9021 - - ksqldb-server: - image: confluentinc/cp-ksqldb-server:7.5.3 - hostname: ksqldb-server - container_name: ksqldb-server - depends_on: - - broker - - connect - ports: - - "8088:8088" - environment: - KSQL_CONFIG_DIR: "/etc/ksql" - KSQL_BOOTSTRAP_SERVERS: "broker:29092" - KSQL_HOST_NAME: ksqldb-server - KSQL_LISTENERS: "http://0.0.0.0:8088" - KSQL_CACHE_MAX_BYTES_BUFFERING: 0 - KSQL_KSQL_SCHEMA_REGISTRY_URL: "http://schema-registry:8081" - KSQL_PRODUCER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringProducerInterceptor" - KSQL_CONSUMER_INTERCEPTOR_CLASSES: "io.confluent.monitoring.clients.interceptor.MonitoringConsumerInterceptor" - KSQL_KSQL_CONNECT_URL: "http://connect:8083" - KSQL_KSQL_LOGGING_PROCESSING_TOPIC_REPLICATION_FACTOR: 1 - KSQL_KSQL_LOGGING_PROCESSING_TOPIC_AUTO_CREATE: 'true' - KSQL_KSQL_LOGGING_PROCESSING_STREAM_AUTO_CREATE: 'true' - - # ksqldb-cli: - # image: confluentinc/cp-ksqldb-cli:7.5.3 - # container_name: ksqldb-cli - # depends_on: - # - broker - # - connect - # - ksqldb-server - # entrypoint: /bin/sh - # tty: true - - # ksql-datagen: - # image: confluentinc/ksqldb-examples:7.5.3 - # hostname: ksql-datagen - # container_name: ksql-datagen - # depends_on: - # - ksqldb-server - # - broker - # - schema-registry - # - connect - # command: "bash -c 'echo Waiting for Kafka to be ready... && \ - # cub kafka-ready -b broker:29092 1 40 && \ - # echo Waiting for Confluent Schema Registry to be ready... && \ - # cub sr-ready schema-registry 8081 40 && \ - # echo Waiting a few seconds for topic creation to finish... && \ - # sleep 11 && \ - # tail -f /dev/null'" - # environment: - # KSQL_CONFIG_DIR: "/etc/ksql" - # STREAMS_BOOTSTRAP_SERVERS: broker:29092 - # STREAMS_SCHEMA_REGISTRY_HOST: schema-registry - # STREAMS_SCHEMA_REGISTRY_PORT: 8081 - - rest-proxy: - image: confluentinc/cp-kafka-rest:7.5.3 + KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 + # Consumer groups are created per-run, so skip the 3s default rebalance debounce + KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 + # Tests create their topics explicitly; test_invalid_topic depends on this being off + KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false" + healthcheck: + test: ["CMD-SHELL", "kafka-broker-api-versions --bootstrap-server localhost:9092"] + interval: 5s + timeout: 10s + retries: 24 depends_on: - - broker - - schema-registry - ports: - - 8082:8082 - hostname: rest-proxy - container_name: rest-proxy - environment: - KAFKA_REST_HOST_NAME: rest-proxy - KAFKA_REST_BOOTSTRAP_SERVERS: 'broker:29092' - KAFKA_REST_LISTENERS: "http://0.0.0.0:8082" - KAFKA_REST_SCHEMA_REGISTRY_URL: 'http://schema-registry:8081' \ No newline at end of file + zoo1: + condition: service_healthy diff --git a/conda/dev-environment-unix.yml b/conda/dev-environment-unix.yml index fadb0b141..9ff9f08e8 100644 --- a/conda/dev-environment-unix.yml +++ b/conda/dev-environment-unix.yml @@ -17,7 +17,6 @@ dependencies: - flex - graphviz - gtest - - httpx>=0.20,<1 - libarrow<24 - libboost>=1.80.0 - libboost-headers>=1.80.0 @@ -41,6 +40,7 @@ dependencies: - pytest-sugar - python>=3.10,<3.15 - python-build + - python-confluent-kafka - python-graphviz - python-rapidjson - pytz diff --git a/conda/dev-environment-win.yml b/conda/dev-environment-win.yml index 65ca65cb4..4ac9e69e4 100644 --- a/conda/dev-environment-win.yml +++ b/conda/dev-environment-win.yml @@ -17,7 +17,6 @@ dependencies: # - flex # not available on windows - graphviz - gtest - - httpx>=0.20,<1 - libarrow<24 - libboost>=1.80.0 - libboost-headers>=1.80.0 @@ -41,6 +40,7 @@ dependencies: - pytest-sugar - python>=3.10,<3.14 - python-build + - python-confluent-kafka - python-graphviz - python-rapidjson - pytz diff --git a/cpp/csp/adapters/kafka/KafkaAdapterManager.cpp b/cpp/csp/adapters/kafka/KafkaAdapterManager.cpp index b6251364d..1bae12fa9 100644 --- a/cpp/csp/adapters/kafka/KafkaAdapterManager.cpp +++ b/cpp/csp/adapters/kafka/KafkaAdapterManager.cpp @@ -23,6 +23,14 @@ INIT_CSP_ENUM( csp::adapters::kafka::KafkaStatusMessageType, namespace csp::adapters::kafka { +//A topic that does not exist is only reported this way when the broker will not auto-create it. +//The producer reports it after topic.metadata.propagation.max.ms rather than immediately. +static bool isUnknownTopic( int err ) +{ + return err == RdKafka::ErrorCode::ERR_UNKNOWN_TOPIC_OR_PART || + err == RdKafka::ErrorCode::ERR__UNKNOWN_TOPIC; +} + class DeliveryReportCb : public RdKafka::DeliveryReportCb { public: @@ -38,7 +46,12 @@ class DeliveryReportCb : public RdKafka::DeliveryReportCb { std::string msg = "KafkaPublisher: Message delivery failed for topic " + message.topic_name() + ". Failure: " + message.errstr(); m_adapterManager -> pushStatus( StatusLevel::ERROR, KafkaStatusMessageType::MSG_DELIVERY_FAILED, msg ); + + if( isUnknownTopic( message.err() ) ) + m_adapterManager -> forceShutdown( msg ); } + else + m_adapterManager -> onBrokerActivity(); } private: KafkaAdapterManager * m_adapterManager; @@ -60,13 +73,17 @@ class EventCb : public RdKafka::EventCb if( event.type() == RdKafka::Event::EVENT_ERROR ) { //We shutdown the app if its a fatal error OR if its an authentication issue which has plagued users multiple times - //Adding ERR__ALL_BROKERS_DOWN which happens when all brokers are down if( event.fatal() || - event.err() == RdKafka::ErrorCode::ERR__AUTHENTICATION || - event.err() == RdKafka::ErrorCode::ERR__ALL_BROKERS_DOWN ) + event.err() == RdKafka::ErrorCode::ERR__AUTHENTICATION ) { m_adapterManager -> forceShutdown( RdKafka::err2str( ( RdKafka::ErrorCode ) event.err() ) + event.str() ); } + //librdkafka reports this on every failed connection round and reconnects on its own, so a + //single report says nothing about whether the brokers are really gone + else if( event.err() == RdKafka::ErrorCode::ERR__ALL_BROKERS_DOWN ) + { + m_adapterManager -> onBrokersDown( RdKafka::err2str( ( RdKafka::ErrorCode ) event.err() ) + event.str() ); + } } } @@ -77,10 +94,12 @@ class EventCb : public RdKafka::EventCb KafkaAdapterManager::KafkaAdapterManager( csp::Engine * engine, const Dictionary & properties ) : AdapterManager( engine ), m_consumerIdx( 0 ), m_producerPollThreadActive( false ), - m_unrecoverableError( false ) + m_unrecoverableError( false ), + m_brokersDownSince( 0 ) { m_maxThreads = properties.get( "max_threads" ); m_pollTimeoutMs = properties.get( "poll_timeout" ).asMilliseconds(); + m_brokersDownTolerance = properties.get( "broker_down_tolerance" ).asNanoseconds(); m_eventCb = std::make_unique( this ); m_producerCb = std::make_unique( this ); @@ -133,6 +152,33 @@ void KafkaAdapterManager::setConfProperties( RdKafka::Conf * conf, const Diction } } +void KafkaAdapterManager::onBrokersDown( const std::string & err ) +{ + int64_t since = 0; + + //Only the first report starts the clock. Escalation is left to checkBrokersDown() on the poll + //threads rather than done here, since librdkafka may not report again for seconds. + if( m_brokersDownSince.compare_exchange_strong( since, DateTime::now().asNanoseconds() ) ) + { + std::lock_guard guard( m_brokersDownLock ); + m_brokersDownError = err; + } +} + +void KafkaAdapterManager::checkBrokersDown() +{ + int64_t since = m_brokersDownSince.load( std::memory_order_relaxed ); + if( since == 0 || DateTime::now().asNanoseconds() - since < m_brokersDownTolerance ) + return; + + std::string err; + { + std::lock_guard guard( m_brokersDownLock ); + err = m_brokersDownError; + } + forceShutdown( err ); +} + void KafkaAdapterManager::forceShutdown( const std::string & err ) { m_unrecoverableError = true; // So we can alert the producer to stop trying to flush @@ -228,6 +274,7 @@ void KafkaAdapterManager::pollProducers() while( m_producerPollThreadActive ) { m_producer -> poll( m_pollTimeoutMs ); + checkBrokersDown(); } try @@ -235,7 +282,9 @@ void KafkaAdapterManager::pollProducers() while( true ) { auto rc = m_producer -> flush( 5000 ); - if( !rc || m_unrecoverableError ) + //Waiting on a flush that cannot complete is what used to hang shutdown, so give up as + //soon as there is nothing to flush to + if( !rc || m_unrecoverableError || brokersDown() ) break; if( rc != RdKafka::ERR__TIMED_OUT ) diff --git a/cpp/csp/adapters/kafka/KafkaAdapterManager.h b/cpp/csp/adapters/kafka/KafkaAdapterManager.h index c095d026a..fbb275125 100644 --- a/cpp/csp/adapters/kafka/KafkaAdapterManager.h +++ b/cpp/csp/adapters/kafka/KafkaAdapterManager.h @@ -8,6 +8,7 @@ #include #include #include +#include #include #include #include @@ -75,6 +76,11 @@ class KafkaAdapterManager final : public csp::AdapterManager void forceShutdown( const std::string & err ); + void onBrokersDown( const std::string & err ); + void checkBrokersDown(); + bool brokersDown() const { return m_brokersDownSince.load( std::memory_order_relaxed ) != 0; } + void onBrokerActivity() { m_brokersDownSince.store( 0, std::memory_order_relaxed ); } + void markConsumerReplayDone( KafkaConsumer * consumer, const std::string & topic ); void onMessage( RdKafka::Message * msg ) const; @@ -132,6 +138,12 @@ class KafkaAdapterManager final : public csp::AdapterManager std::atomic m_producerPollThreadActive; std::atomic m_unrecoverableError; + //Nanoseconds since the brokers were first reported down, 0 while they are believed up + std::atomic m_brokersDownSince; + int64_t m_brokersDownTolerance; + std::mutex m_brokersDownLock; + std::string m_brokersDownError; + std::unique_ptr m_consumerConf; std::unique_ptr m_producerConf; Dictionary::Value m_startOffsetProperty; diff --git a/cpp/csp/adapters/kafka/KafkaConsumer.cpp b/cpp/csp/adapters/kafka/KafkaConsumer.cpp index de353e639..2031d12cf 100644 --- a/cpp/csp/adapters/kafka/KafkaConsumer.cpp +++ b/cpp/csp/adapters/kafka/KafkaConsumer.cpp @@ -196,6 +196,8 @@ void KafkaConsumer::poll() { std::unique_ptr msg( m_consumer -> consume( m_mgr -> pollTimeoutMs() ) ); + m_mgr -> checkBrokersDown(); + if( msg -> err() == RdKafka::ERR__TIMED_OUT ) continue; @@ -208,8 +210,20 @@ void KafkaConsumer::poll() continue; } + //Only reported when the broker will not auto-create the topic. Waiting for a topic + //that will never appear looks identical to an idle subscription, so surface it. + if( msg -> err() == RdKafka::ERR_UNKNOWN_TOPIC_OR_PART ) [[unlikely]] + { + m_mgr -> forceShutdown( RdKafka::err2str( msg -> err() ) + " error: " + msg -> errstr() ); + continue; + } + if( msg -> err() == RdKafka::ERR_NO_ERROR && msg -> len() ) + { + //Proof the brokers came back, whatever was reported down before + m_mgr -> onBrokerActivity(); m_mgr -> onMessage( msg.get() ); + } //Not sure why, but it looks like we repeatedly get EOF callbacks even after the original one //may want to look into this. Not an issue in practice, but seems like unnecessary overhead else if( msg -> err() == RdKafka::ERR__PARTITION_EOF ) diff --git a/cpp/csp/engine/Struct.h b/cpp/csp/engine/Struct.h index 7082fefc9..f4c4f5958 100644 --- a/cpp/csp/engine/Struct.h +++ b/cpp/csp/engine/Struct.h @@ -737,7 +737,7 @@ std::shared_ptr::type> StructMeta::getMetaField( std::shared_ptr::type> typedfield = std::dynamic_pointer_cast::type>( field_ ); if( !typedfield ) CSP_THROW( TypeError, expectedtype << " - provided struct type " << name() << " expected type " << CspType::Type::fromCType::type << " for field " << fieldname - << " but got type " << field_ -> type() -> type() << " for " << expectedtype ); + << " but got type " << field_ -> type() -> type() ); return typedfield; } diff --git a/csp/adapters/kafka.py b/csp/adapters/kafka.py index daa3282d0..0dc368838 100644 --- a/csp/adapters/kafka.py +++ b/csp/adapters/kafka.py @@ -20,7 +20,18 @@ from csp.impl.wiring import ReplayMode, input_adapter_def, output_adapter_def, status_adapter_def from csp.lib import _kafkaadapterimpl -_ = BytesMessageProtoMapper, DateTimeType, JSONTextMessageMapper, RawBytesMessageMapper, RawTextMessageMapper +__all__ = ( + "KafkaAdapterManager", + "KafkaStartOffset", + "KafkaStatusMessageType", + # re-exported from csp.adapters.utils for backwards compatibility + "BytesMessageProtoMapper", + "DateTimeType", + "JSONTextMessageMapper", + "RawBytesMessageMapper", + "RawTextMessageMapper", +) + T = TypeVar("T") @@ -54,6 +65,7 @@ def __init__( rd_kafka_conf_options=None, debug: bool = False, poll_timeout: timedelta = timedelta(seconds=1), + broker_down_tolerance: timedelta = timedelta(milliseconds=500), rd_kafka_consumer_conf_options=None, rd_kafka_producer_conf_options=None, ): @@ -67,6 +79,8 @@ def __init__( set in this case since adapter will always replay from the last consumed offset. :param group_id_prefix - ( optional ) when not passing an explicit group_id, a prefix can be supplied that will be use to prefix the UUID generated for the group_id + :param broker_down_tolerance - how long all brokers may stay unreachable before the engine is shut down. librdkafka + reconnects on its own, so a single report is not evidence the brokers are gone. """ if group_id is not None and start_offset is not None: raise ValueError("start_offset is not supported when consuming with group_id") @@ -101,6 +115,7 @@ def __init__( "start_offset": start_offset.value if isinstance(start_offset, KafkaStartOffset) else start_offset, "max_threads": max_threads, "poll_timeout": poll_timeout, + "broker_down_tolerance": broker_down_tolerance, "rd_kafka_conf_properties": conf_properties, "rd_kafka_consumer_conf_properties": consumer_properties, "rd_kafka_producer_conf_properties": producer_properties, diff --git a/csp/tests/adapters/conftest.py b/csp/tests/adapters/conftest.py index 62046fa4a..229fd4626 100644 --- a/csp/tests/adapters/conftest.py +++ b/csp/tests/adapters/conftest.py @@ -1,3 +1,5 @@ +from uuid import uuid4 + import pytest from csp.adapters.kafka import KafkaAdapterManager @@ -5,15 +7,17 @@ @pytest.fixture(scope="module", autouse=True) def kafkabroker(): - return "localhost:9092" + # Defined in ci/kafka/docker-compose.yml. Not "localhost": that also resolves to ::1, where + # the container publishes nothing, and librdkafka alternates between the resolved addresses. + return "127.0.0.1:9092" @pytest.fixture(scope="module", autouse=True) def kafkaadapterkwargs(kafkabroker): - return dict(broker=kafkabroker, group_id="group.id123", rd_kafka_conf_options={"allow.auto.create.topics": "true"}) + # Unique group id so a rerun never inherits committed offsets from a previous run + return dict(broker=kafkabroker, group_id=f"csp.test.{uuid4()}") @pytest.fixture(scope="module", autouse=True) def kafkaadapter(kafkaadapterkwargs): - _kafkaadapter = KafkaAdapterManager(**kafkaadapterkwargs) - return _kafkaadapter + return KafkaAdapterManager(**kafkaadapterkwargs) diff --git a/csp/tests/adapters/kafka_utils.py b/csp/tests/adapters/kafka_utils.py index 98544cd34..42ccbd3e1 100644 --- a/csp/tests/adapters/kafka_utils.py +++ b/csp/tests/adapters/kafka_utils.py @@ -1,13 +1,16 @@ -import httpx - - -def _precreate_topic(topic): - """Since we test against confluent kafka, just use the kafka rest addon""" - rest_broker = "http://localhost:8082" - cluster_info = httpx.get(f"{rest_broker}/v3/clusters") - cluster_id = cluster_info.json()["data"][0]["cluster_id"] - resp = httpx.post(f"{rest_broker}/v3/clusters/{cluster_id}/topics", json={"topic_name": topic}) - if resp.status_code != 201 and "already exists" not in resp.content.decode("utf8"): - raise Exception( - f"Could not create topic {topic} on cluster {rest_broker}/v3/clusters/{cluster_id}/topics - received {resp.content}" - ) +__all__ = ("create_topic",) + + +def create_topic(broker, topic): + """Create `topic` and block until the broker acknowledges it. + + Creation is done out of band rather than by publishing a warm-up message, so that no test data + lands on the topic and so tests do not depend on broker-side auto-creation (disabled in + ci/kafka/docker-compose.yml so test_invalid_topic can exercise the failure path). + """ + # Imported lazily so collection does not require confluent-kafka when the kafka tests are skipped + from confluent_kafka.admin import AdminClient, NewTopic + + admin = AdminClient({"bootstrap.servers": broker}) + for _, future in admin.create_topics([NewTopic(topic, num_partitions=1, replication_factor=1)]).items(): + future.result() diff --git a/csp/tests/adapters/test_kafka.py b/csp/tests/adapters/test_kafka.py index 6cc15388f..0724d2162 100644 --- a/csp/tests/adapters/test_kafka.py +++ b/csp/tests/adapters/test_kafka.py @@ -6,17 +6,11 @@ import csp from csp import ts -from csp.adapters.kafka import ( - DateTimeType, - JSONTextMessageMapper, - KafkaAdapterManager, - KafkaStartOffset, - RawBytesMessageMapper, - RawTextMessageMapper, -) +from csp.adapters.kafka import KafkaAdapterManager, KafkaStartOffset +from csp.adapters.utils import DateTimeType, JSONTextMessageMapper, RawBytesMessageMapper, RawTextMessageMapper from csp.utils.datetime import utc_now -from .kafka_utils import _precreate_topic +from .kafka_utils import create_topic class MyData(csp.Struct): @@ -69,7 +63,10 @@ class MetaSubData(csp.Struct): class TestKafka: @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") - def test_metadata(self, kafkaadapter): + def test_metadata(self, kafkaadapter, kafkabroker): + topic = f"test.metadata.{os.getpid()}" + create_topic(kafkabroker, topic) + def graph(count: int): msg_mapper = JSONTextMessageMapper(datetime_type=DateTimeType.UINT64_MICROS) @@ -84,19 +81,17 @@ def graph(count: int): "timestamp": "mapped_timestamp", } - topic = f"test.metadata.{os.getpid()}" - _precreate_topic(topic) subKey = "foo" pubKey = ["mapped_a", "mapped_b", "mapped_c"] - c = csp.count(csp.timer(timedelta(seconds=0.1))) + # Publish slowly enough that the consumer reaches end-of-partition and flips to live + c = csp.count(csp.timer(timedelta(seconds=1))) t = csp.sample(c, csp.const("foo")) pubStruct = MetaPubData.collectts( mapped_a=MetaSubStruct.collectts(mapped_b=MetaTextStruct.collectts(mapped_c=t)), mapped_count=c ) - # csp.print('pub', pubStruct) kafkaadapter.publish(msg_mapper, topic, pubKey, pubStruct, field_map=pub_field_map) sub_data = kafkaadapter.subscribe( @@ -110,7 +105,7 @@ def graph(count: int): ) csp.add_graph_output("sub_data", sub_data) - # csp.print('sub', sub_data) + # Wait for at least count ticks and until we get a live tick done_flag = csp.count(sub_data) >= count done_flag = csp.and_(done_flag, sub_data.mapped_live == True) # noqa: E712 @@ -119,17 +114,18 @@ def graph(count: int): count = 5 results = csp.run(graph, count, starttime=utc_now(), endtime=timedelta(seconds=30), realtime=True) - assert len(results["sub_data"]) >= 5 - print(results) + assert len(results["sub_data"]) >= count + for result in results["sub_data"]: assert result[1].mapped_partition >= 0 assert result[1].mapped_offset >= 0 assert result[1].mapped_live is not None assert result[1].mapped_timestamp < utc_now() + # last record should be live (first may or may not be live depending on timing) assert results["sub_data"][-1][1].mapped_live @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") - def test_basic(self, kafkaadapter): + def test_basic(self, kafkaadapter, kafkabroker): @csp.node def curtime(x: ts[object]) -> ts[datetime]: if csp.ticked(x): @@ -140,6 +136,7 @@ def graph(symbols: list, count: int): csp.timer(timedelta(seconds=0.2), True), csp.delay(csp.timer(timedelta(seconds=0.2), False), timedelta(seconds=0.1)), ) + i = csp.count(csp.timer(timedelta(seconds=0.15))) d = csp.count(csp.timer(timedelta(seconds=0.2))) / 2.0 s = csp.sample(csp.timer(timedelta(seconds=0.4)), csp.const("STRING")) @@ -153,8 +150,7 @@ def graph(symbols: list, count: int): struct_field_map = {"b": "b2", "i": "i2", "d": "d2", "s": "s2", "dt": "dt2", "date": "date2"} done_flags = [] - topic = f"mktdata.{os.getpid()}" - _precreate_topic(topic) + for symbol in symbols: kafkaadapter.publish(msg_mapper, topic, symbol, b, field_map="b") kafkaadapter.publish(msg_mapper, topic, symbol, i, field_map="i") @@ -181,8 +177,6 @@ def graph(symbols: list, count: int): ) csp.add_graph_output(f"pall_{symbol}", pub_data) - # csp.print('status', kafkaadapter.status()) - sub_data = kafkaadapter.subscribe( ts_type=SubData, msg_mapper=msg_mapper, @@ -190,9 +184,6 @@ def graph(symbols: list, count: int): key=symbol, push_mode=csp.PushMode.NON_COLLAPSING, ) - - sub_data = csp.firstN(sub_data, count) - csp.add_graph_output(f"sall_{symbol}", sub_data) done_flag = csp.count(sub_data) == count @@ -203,20 +194,25 @@ def graph(symbols: list, count: int): stop = csp.filter(stop, stop) csp.stop_engine(stop) + topic = f"mktdata.{os.getpid()}" + create_topic(kafkabroker, topic) symbols = ["AAPL", "MSFT"] - count = 100 + count = 50 results = csp.run(graph, symbols, count, starttime=utc_now(), endtime=timedelta(seconds=30), realtime=True) for symbol in symbols: - pub = results[f"pall_{symbol}"] - sub = results[f"sall_{symbol}"] + pub = [v[1] for v in results[f"pall_{symbol}"]] + sub = [v[1] for v in results[f"sall_{symbol}"]] + # The subscriber may miss a prefix of the stream while its consumer group is being + # assigned, so align on the first record it did see rather than assuming it saw all. assert len(sub) == count - assert [v[1] for v in sub] == [v[1] for v in pub[:count]] + start = pub.index(sub[0]) + assert sub == pub[start : start + count] @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") def test_start_offsets(self, kafkaadapter, kafkabroker): topic = f"test_start_offsets.{os.getpid()}" - _precreate_topic(topic) + create_topic(kafkabroker, topic) msg_mapper = JSONTextMessageMapper(datetime_type=DateTimeType.UINT64_MICROS) count = 10 @@ -228,7 +224,6 @@ def pub_graph(): stop = csp.count(struct) == count stop = csp.filter(stop, stop) csp.stop_engine(stop) - # csp.print('pub', struct) csp.run(pub_graph, starttime=utc_now(), endtime=timedelta(seconds=30), realtime=True) @@ -247,9 +242,6 @@ def get_times_graph(): csp.stop_engine(csp.filter(stop, stop)) csp.add_graph_output("data", data) - # csp.print('sub', data) - # csp.print('status', kafkaadapter.status()) - all_data = csp.run(get_times_graph, starttime=utc_now(), endtime=timedelta(seconds=30), realtime=True)["data"] min_time = all_data[0][1].dt @@ -267,8 +259,6 @@ def get_data(start_offset, expected_count): csp.stop_engine(csp.filter(stop, stop)) csp.add_graph_output("data", data) - # csp.print('data', data) - res = csp.run( get_data, KafkaStartOffset.EARLIEST, @@ -277,7 +267,6 @@ def get_data(start_offset, expected_count): endtime=timedelta(seconds=30), realtime=True, )["data"] - # print(res) # If we playback from earliest but start "now", all data should still arrive but as realtime ticks assert len(res) == 10 @@ -292,7 +281,7 @@ def get_data(start_offset, expected_count): assert len(res) == 0 res = csp.run( - get_data, KafkaStartOffset.START_TIME, 10, starttime=min_time, endtime=timedelta(seconds=30), realtime=True + get_data, KafkaStartOffset.START_TIME, 10, starttime=min_time, endtime=timedelta(seconds=10), realtime=True )["data"] assert len(res) == 10 @@ -308,12 +297,12 @@ def get_data(start_offset, expected_count): assert len(res) == len(expected) res = csp.run( - get_data, timedelta(seconds=0), len(expected), starttime=stime, endtime=timedelta(seconds=30), realtime=True + get_data, timedelta(seconds=0), len(expected), starttime=stime, endtime=timedelta(seconds=10), realtime=True )["data"] assert len(res) == len(expected) @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") - def test_raw_pubsub(self, kafkaadapter): + def test_raw_pubsub(self, kafkaadapter, kafkabroker): @csp.node def data(x: ts[object]) -> ts[bytes]: if csp.ticked(x): @@ -329,15 +318,10 @@ def graph(symbols: list, count: int): msg_mapper = RawBytesMessageMapper() done_flags = [] - topic = f"test_str.{os.getpid()}" - _precreate_topic(topic) for symbol in symbols: - topic = f"test_str.{os.getpid()}" kafkaadapter.publish(msg_mapper, topic, symbol, d) csp.add_graph_output(f"pub_{symbol}", d) - # csp.print('status', kafkaadapter.status()) - sub_data = kafkaadapter.subscribe( ts_type=SubData, msg_mapper=RawTextMessageMapper(), @@ -356,14 +340,12 @@ def graph(symbols: list, count: int): push_mode=csp.PushMode.NON_COLLAPSING, ) - sub_data = csp.firstN(sub_data.msg, count) - sub_data_bytes = csp.firstN(sub_data_bytes, count) - - # csp.print('sub', sub_data) - csp.add_graph_output(f"sub_{symbol}", sub_data) + csp.add_graph_output(f"sub_{symbol}", sub_data.msg) csp.add_graph_output(f"sub_bytes_{symbol}", sub_data_bytes) - done_flag = csp.count(sub_data) + csp.count(sub_data_bytes) == count * 2 + # Wait for count messages on both subscribers + done_flag = csp.count(sub_data) >= count + done_flag = csp.and_(done_flag, csp.count(sub_data_bytes) >= count) done_flag = csp.filter(done_flag, done_flag) done_flags.append(done_flag) @@ -371,29 +353,35 @@ def graph(symbols: list, count: int): stop = csp.filter(stop, stop) csp.stop_engine(stop) + topic = f"test_str.{os.getpid()}" + create_topic(kafkabroker, topic) + symbols = ["AAPL", "MSFT"] count = 10 results = csp.run(graph, symbols, count, starttime=utc_now(), endtime=timedelta(seconds=30), realtime=True) - # print(results) for symbol in symbols: - pub = results[f"pub_{symbol}"] - sub = results[f"sub_{symbol}"] - sub_bytes = results[f"sub_bytes_{symbol}"] - - assert len(sub) == count - assert [v[1] for v in sub] == [v[1] for v in pub[:count]] - assert [v[1] for v in sub_bytes] == [v[1] for v in pub[:count]] + pub = [v[1] for v in results[f"pub_{symbol}"]] + sub = [v[1] for v in results[f"sub_{symbol}"]] + sub_bytes = [v[1] for v in results[f"sub_bytes_{symbol}"]] + + # Align on the first record each subscriber saw, then require an exact contiguous run. + # Comparing sequences (rather than set membership) is what catches loss, duplication + # and reordering. + for received in (sub, sub_bytes): + assert len(received) >= count + start = pub.index(received[0]) + assert received == pub[start : start + len(received)] @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") def test_invalid_topic(self, kafkaadapterkwargs): class SubData(csp.Struct): msg: str + # Relies on broker-side auto.create.topics.enable=false (ci/kafka/docker-compose.yml) kafkaadapter1 = KafkaAdapterManager(**kafkaadapterkwargs) # Was a bug where engine would stall def graph_sub(): - # csp.print('status', kafkaadapter.status()) return kafkaadapter1.subscribe( ts_type=SubData, msg_mapper=RawTextMessageMapper(), field_map={"": "msg"}, topic="foobar", key="none" ) @@ -401,7 +389,13 @@ def graph_sub(): # With bug this would deadlock with pytest.raises(RuntimeError): csp.run(graph_sub, starttime=utc_now(), endtime=timedelta(seconds=2), realtime=True) - kafkaadapter2 = KafkaAdapterManager(**kafkaadapterkwargs) + + kafkaadapter2 = KafkaAdapterManager( + **kafkaadapterkwargs, + # The producer holds a message for a topic it cannot find until this elapses, on the + # assumption the topic is about to be created; the 30s default dominates the test + rd_kafka_producer_conf_options={"topic.metadata.propagation.max.ms": "2000"}, + ) def graph_pub(): msg_mapper = RawTextMessageMapper() @@ -409,7 +403,7 @@ def graph_pub(): # With bug this would deadlock with pytest.raises(RuntimeError): - csp.run(graph_pub, starttime=utc_now(), endtime=timedelta(seconds=2), realtime=True) + csp.run(graph_pub, starttime=utc_now(), endtime=timedelta(seconds=5), realtime=True) @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") def test_invalid_broker(self, kafkaadapterkwargs): @@ -442,15 +436,13 @@ def graph_pub(): csp.run(graph_pub, starttime=utc_now(), endtime=timedelta(seconds=2), realtime=True) @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") - def test_meta_field_map_tick_timestamp_from_field(self, kafkaadapterkwargs): + def test_meta_field_map_tick_timestamp_from_field(self, kafkaadapter): class SubData(csp.Struct): msg: str dt: datetime - kafkaadapter1 = KafkaAdapterManager(**kafkaadapterkwargs) - def graph_sub(): - return kafkaadapter1.subscribe( + return kafkaadapter.subscribe( ts_type=SubData, msg_mapper=RawTextMessageMapper(), meta_field_map={"timestamp": "dt"}, @@ -482,7 +474,7 @@ class BasicData(csp.Struct): b: bool topic = f"test_burst.{os.getpid()}" - _precreate_topic(topic) + create_topic(kafkabroker, topic) msg_mapper = JSONTextMessageMapper(datetime_type=DateTimeType.UINT64_MICROS) count = 10 diff --git a/csp/tests/adapters/test_status.py b/csp/tests/adapters/test_status.py index 9ad5e7147..0977478cb 100644 --- a/csp/tests/adapters/test_status.py +++ b/csp/tests/adapters/test_status.py @@ -1,24 +1,25 @@ import os -from datetime import datetime, timedelta +from datetime import timedelta import pytest import csp from csp import ts -from csp.adapters.kafka import DateTimeType, JSONTextMessageMapper, KafkaStatusMessageType +from csp.adapters.kafka import KafkaStatusMessageType from csp.adapters.status import Level +from csp.adapters.utils import DateTimeType, JSONTextMessageMapper from csp.utils.datetime import utc_now -from .kafka_utils import _precreate_topic +from .kafka_utils import create_topic class SubData(csp.Struct): a: bool -class TestStatus: +class TestStatusKafka: @pytest.mark.skipif(not os.environ.get("CSP_TEST_KAFKA"), reason="Skipping kafka adapter tests") - def test_basic(self, kafkaadapter): + def test_basic(self, kafkaadapter, kafkabroker): topic = f"csp.unittest.{os.getpid()}" key = "test_status" @@ -43,7 +44,7 @@ def graph(): done_flag = csp.count(status) == 1 csp.stop_engine(done_flag) - _precreate_topic(topic) + create_topic(kafkabroker, topic) results = csp.run(graph, starttime=utc_now(), endtime=timedelta(seconds=10), realtime=True) status = results["status"][0][1] assert status.status_code == KafkaStatusMessageType.MSG_RECV_ERROR diff --git a/examples/03_using_adapters/kafka/e1_kafka.py b/examples/03_using_adapters/kafka/e1_kafka.py index a19ace9ca..4ac1eb2e9 100644 --- a/examples/03_using_adapters/kafka/e1_kafka.py +++ b/examples/03_using_adapters/kafka/e1_kafka.py @@ -4,10 +4,10 @@ import csp from csp import ts from csp.adapters.kafka import ( + BytesMessageProtoMapper, DateTimeType, JSONTextMessageMapper, KafkaAdapterManager, - ProtoMessageMapper, RawTextMessageMapper, ) from csp.utils.datetime import utc_now @@ -142,7 +142,7 @@ def proto_graph(): "px": "price", } - msg_mapper = ProtoMessageMapper( + msg_mapper = BytesMessageProtoMapper( proto_directory="/tmp", proto_filename="fxspotstream.proto", proto_message="Snapshot" ) @@ -157,7 +157,7 @@ def proto_graph_multiple_subscribers(): topic = "test2" - msg_mapper = ProtoMessageMapper( + msg_mapper = BytesMessageProtoMapper( proto_directory="/tmp", proto_filename="fxspotstream.proto", proto_message="Snapshot" ) diff --git a/pyproject.toml b/pyproject.toml index 96396ee30..de892a3ed 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -81,9 +81,10 @@ develop = [ "graphviz", "pillow", # adapters - "httpx>=0.20,<1", # kafka + "confluent-kafka", # kafka test topic setup "perspective-python>=2", # perspective - "anywidget", # perspective + "anywidget", # perspective (perspective-python 3.x) + "ipywidgets", # perspective (csp.impl.pandas_perspective) "polars", # parquet "psutil", # test_engine/test_history "sqlalchemy", # db @@ -96,21 +97,25 @@ showgraph = [ "pillow", ] test = [ - "graphviz", - "pillow", "pytest", "pytest-asyncio", "pytest-cov", "pytest-sugar", - "httpx>=0.20,<1", - "perspective-python", - "polars", - "psutil", - "requests", - "slack-sdk>=3", - "sqlalchemy", - "threadpoolctl", - "tornado", + # showgraph + "graphviz", + "pillow", + # adapters + "confluent-kafka", # kafka test topic setup + "perspective-python>=2", # perspective + "anywidget", # perspective (perspective-python 3.x) + "ipywidgets", # perspective (csp.impl.pandas_perspective) + "polars", # parquet + "psutil", # test_engine/test_history + "sqlalchemy", # db + "threadpoolctl", # test_random + "requests", # profiler + "tornado", # profiler, perspective, websocket + "python-rapidjson", # websocket ] symphony = [ "csp-adapter-symphony",