6from .
import classes, service
9def pytest_addoption(parser):
10 group = parser.getgroup(
'kafka')
11 group.addoption(
'--kafka')
14 help=
'Disable use of Kafka',
19def pytest_configure(config):
20 config.addinivalue_line(
22 'kafka: per-test Kafka initialization',
26def pytest_service_register(register_service):
27 register_service(
'kafka', service.create_kafka_service)
30@pytest.fixture(scope='session')
31async def _kafka_global_producer(
34) -> typing.AsyncGenerator[classes.KafkaProducer, None]:
35 producer = classes.KafkaProducer(
36 enabled=_kafka_service,
37 bootstrap_servers=_bootstrap_servers,
39 await producer.start()
41 async with contextlib.aclosing(producer):
46async def kafka_producer(
47 _kafka_global_producer,
48) -> typing.AsyncGenerator[classes.KafkaProducer, None]:
50 Per test Kafka producer instance.
52 @returns @ref testsuite.databases.kafka.classes.KafkaProducer
54 @ingroup userver_testsuite_fixtures
55 Part of the [yandex-taxi-testsuite](https://github.com/yandex/yandex-taxi-testsuite/blob/develop/testsuite/databases/kafka/pytest_plugin.py#L46)
57 yield _kafka_global_producer
58 await _kafka_global_producer._flush()
61@pytest.fixture(scope='session')
62async def _kafka_global_consumer(
65) -> typing.AsyncGenerator[classes.KafkaConsumer, None]:
66 consumer = classes.KafkaConsumer(
67 enabled=_kafka_service,
68 bootstrap_servers=_bootstrap_servers,
70 await consumer.start()
72 async with contextlib.aclosing(consumer):
77async def kafka_consumer(
78 _kafka_global_consumer,
79) -> typing.AsyncGenerator[classes.KafkaConsumer, None]:
81 Per test Kafka consumer instance.
83 @returns @ref testsuite.databases.kafka.classes.KafkaConsumer
85 @ingroup userver_testsuite_fixtures
86 Part of the [yandex-taxi-testsuite](https://github.com/yandex/yandex-taxi-testsuite/blob/develop/testsuite/databases/kafka/pytest_plugin.py#L74)
88 yield _kafka_global_consumer
89 await _kafka_global_consumer._unsubscribe()
92@pytest.fixture(scope='session')
93def kafka_custom_topics() -> dict[str, int]:
95 Redefine this fixture to pass your custom dictionary of topics' settings.
97 @ingroup userver_testsuite_fixtures
98 Part of the [yandex-taxi-testsuite](https://github.com/yandex/yandex-taxi-testsuite/blob/develop/testsuite/databases/kafka/pytest_plugin.py#L87)
101 return service.try_get_custom_topics()
104@pytest.fixture(scope='session')
105def kafka_local() -> classes.BootstrapServers:
107 Override to use custom local cluster bootstrap servers.
108 If not empty, no service started.
110 @ingroup userver_testsuite_fixtures
111 Part of the [yandex-taxi-testsuite](https://github.com/yandex/yandex-taxi-testsuite/blob/develop/testsuite/databases/kafka/pytest_plugin.py#L96)
117@pytest.fixture(scope='session')
118def kafka_disabled(pytestconfig) -> bool:
119 return pytestconfig.option.no_kafka
122@pytest.fixture(scope='session')
123def _kafka_service_settings(kafka_custom_topics) -> classes.ServiceSettings:
124 return service.get_service_settings(kafka_custom_topics)
127@pytest.fixture(scope='session')
128def _bootstrap_servers(kafka_local, _kafka_service_settings) -> str:
130 return ','.join(kafka_local)
132 server_host = _kafka_service_settings.server_host
133 server_port = _kafka_service_settings.server_port
134 return f
'{server_host}:{server_port}'
137@pytest.fixture(scope='session')
139 ensure_service_started,
143 _kafka_service_settings,
147 if not kafka_local
and not pytestconfig.option.kafka:
148 ensure_service_started(
'kafka', settings=_kafka_service_settings)