|
1 | | -"""Metrics reporting, for the sync and async Unleash clients.""" |
2 | | - |
3 | | -from typing import Optional |
| 1 | +"""Metrics reporting, for the sync Unleash client.""" |
4 | 2 |
|
5 | 3 | from yggdrasil_engine.engine import UnleashEngine |
6 | 4 |
|
7 | | -from UnleashClient._async_scheduler import _AsyncJob, _AsyncScheduler |
8 | | -from UnleashClient._async_transport import _AsyncTransport |
9 | 5 | from UnleashClient._payloads import _build_metrics_payload |
10 | 6 | from UnleashClient._scheduler import _ScheduledJob, _Scheduler |
11 | 7 | from UnleashClient._transport import _Transport |
@@ -98,87 +94,3 @@ def stop(self) -> None: |
98 | 94 | self.flush() |
99 | 95 | self._scheduler.cancel(self._job) |
100 | 96 | self._job = None |
101 | | - |
102 | | - |
103 | | -class _AsyncMetricsReporter: |
104 | | - """ |
105 | | - Sends feature and impact metrics to Unleash on a recurring interval. |
106 | | -
|
107 | | - :meth:`start` must be called, and :meth:`stop` awaited, from the event loop the client |
108 | | - runs on, and the loop must stay open for as long as metrics are being reported. |
109 | | -
|
110 | | - Example:: |
111 | | -
|
112 | | - reporter = _AsyncMetricsReporter( |
113 | | - config=config, |
114 | | - transport=transport, |
115 | | - scheduler=scheduler, |
116 | | - engine=engine, |
117 | | - impact_metrics=impact_metrics, |
118 | | - ) |
119 | | - reporter.start() |
120 | | -
|
121 | | - await reporter.flush() |
122 | | -
|
123 | | - await reporter.stop() |
124 | | - """ |
125 | | - |
126 | | - def __init__( |
127 | | - self, |
128 | | - config: UnleashConfig, |
129 | | - transport: _AsyncTransport, |
130 | | - scheduler: _AsyncScheduler, |
131 | | - engine: UnleashEngine, |
132 | | - impact_metrics: ImpactMetrics, |
133 | | - ) -> None: |
134 | | - self._config: UnleashConfig = config |
135 | | - self._transport: _AsyncTransport = transport |
136 | | - self._scheduler: _AsyncScheduler = scheduler |
137 | | - self._engine: UnleashEngine = engine |
138 | | - self._impact_metrics: ImpactMetrics = impact_metrics |
139 | | - self._job: Optional[_AsyncJob] = None |
140 | | - |
141 | | - def start(self) -> None: |
142 | | - """Schedules a send every ``metrics_interval`` seconds, with ``metrics_jitter`` of jitter.""" |
143 | | - self._job = self._scheduler.every( |
144 | | - interval_seconds=int(self._config.metrics_interval), |
145 | | - jitter_seconds=self._config.metrics_jitter, |
146 | | - fn=self.flush, |
147 | | - ) |
148 | | - |
149 | | - async def flush(self) -> None: |
150 | | - """ |
151 | | - Sends one bucket of feature and impact metrics. |
152 | | -
|
153 | | - Sends nothing when neither has anything to report. When a send fails or is |
154 | | - cancelled, its impact metrics are restored so the next send carries them. |
155 | | - """ |
156 | | - bucket = self._engine.get_metrics() |
157 | | - impact_metrics = self._impact_metrics.collect() |
158 | | - |
159 | | - if not (bucket or impact_metrics): |
160 | | - LOGGER.debug("No feature flags with metrics, skipping metrics submission.") |
161 | | - return |
162 | | - |
163 | | - payload = _build_metrics_payload(self._config, bucket, impact_metrics) |
164 | | - sent = False |
165 | | - try: |
166 | | - sent = await self._transport.send_metrics(payload) |
167 | | - finally: |
168 | | - if not sent and impact_metrics: |
169 | | - self._impact_metrics.restore(impact_metrics) |
170 | | - |
171 | | - async def stop(self) -> None: |
172 | | - """ |
173 | | - Stops the recurring send and flushes whatever is left. |
174 | | -
|
175 | | - Does nothing when :meth:`start` was never called. A send still in flight is |
176 | | - cancelled: its impact metrics go out with the final flush, and its feature |
177 | | - metrics are lost. |
178 | | - """ |
179 | | - if self._job is None: |
180 | | - return |
181 | | - |
182 | | - job, self._job = self._job, None |
183 | | - await self._scheduler.cancel_and_wait(job) |
184 | | - await self.flush() |
0 commit comments