userver: en/testsuite/databases/kafka/pytest_plugin.py Source File
Loading...
Searching...
No Matches
pytest_plugin.py
1import contextlib
2import typing
3
4import pytest
5
6from . import classes, service
7
8
9def pytest_addoption(parser):
10 group = parser.getgroup('kafka')
11 group.addoption('--kafka')
12 group.addoption(
13 '--no-kafka',
14 help='Disable use of Kafka',
15 action='store_true',
16 )
17
18
19def pytest_configure(config):
20 config.addinivalue_line(
21 'markers',
22 'kafka: per-test Kafka initialization',
23 )
24
25
26def pytest_service_register(register_service):
27 register_service('kafka', service.create_kafka_service)
28
29
30@pytest.fixture(scope='session')
31async def _kafka_global_producer(
32 _kafka_service,
33 _bootstrap_servers,
34) -> typing.AsyncGenerator[classes.KafkaProducer, None]:
35 producer = classes.KafkaProducer(
36 enabled=_kafka_service,
37 bootstrap_servers=_bootstrap_servers,
38 )
39 await producer.start()
40
41 async with contextlib.aclosing(producer):
42 yield producer
43
44
45@pytest.fixture
46async def kafka_producer(
47 _kafka_global_producer,
48) -> typing.AsyncGenerator[classes.KafkaProducer, None]:
49 """
50 Per test Kafka producer instance.
51
52 @returns @ref testsuite.databases.kafka.classes.KafkaProducer
53
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)
56 """
57 yield _kafka_global_producer
58 await _kafka_global_producer._flush()
59
60
61@pytest.fixture(scope='session')
62async def _kafka_global_consumer(
63 _kafka_service,
64 _bootstrap_servers,
65) -> typing.AsyncGenerator[classes.KafkaConsumer, None]:
66 consumer = classes.KafkaConsumer(
67 enabled=_kafka_service,
68 bootstrap_servers=_bootstrap_servers,
69 )
70 await consumer.start()
71
72 async with contextlib.aclosing(consumer):
73 yield consumer
74
75
76@pytest.fixture
77async def kafka_consumer(
78 _kafka_global_consumer,
79) -> typing.AsyncGenerator[classes.KafkaConsumer, None]:
80 """
81 Per test Kafka consumer instance.
82
83 @returns @ref testsuite.databases.kafka.classes.KafkaConsumer
84
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)
87 """
88 yield _kafka_global_consumer
89 await _kafka_global_consumer._unsubscribe()
90
91
92@pytest.fixture(scope='session')
93def kafka_custom_topics() -> dict[str, int]:
94 """
95 Redefine this fixture to pass your custom dictionary of topics' settings.
96
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)
99 """
100
101 return service.try_get_custom_topics()
102
103
104@pytest.fixture(scope='session')
105def kafka_local() -> classes.BootstrapServers:
106 """
107 Override to use custom local cluster bootstrap servers.
108 If not empty, no service started.
109
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)
112 """
113
114 return []
115
116
117@pytest.fixture(scope='session')
118def kafka_disabled(pytestconfig) -> bool:
119 return pytestconfig.option.no_kafka
120
121
122@pytest.fixture(scope='session')
123def _kafka_service_settings(kafka_custom_topics) -> classes.ServiceSettings:
124 return service.get_service_settings(kafka_custom_topics)
125
126
127@pytest.fixture(scope='session')
128def _bootstrap_servers(kafka_local, _kafka_service_settings) -> str:
129 if kafka_local:
130 return ','.join(kafka_local)
131
132 server_host = _kafka_service_settings.server_host
133 server_port = _kafka_service_settings.server_port
134 return f'{server_host}:{server_port}'
135
136
137@pytest.fixture(scope='session')
138def _kafka_service(
139 ensure_service_started,
140 kafka_local,
141 kafka_disabled,
142 pytestconfig,
143 _kafka_service_settings,
144) -> bool:
145 if kafka_disabled:
146 return False
147 if not kafka_local and not pytestconfig.option.kafka:
148 ensure_service_started('kafka', settings=_kafka_service_settings)
149 return True