session
session
¶
Functions:
| Name | Description |
|---|---|
measurement_session |
Activate a collector for code running in this context. |
configured_measurement_session |
Activate and persist a collector when a measurement config is provided. |
current_collector |
Return the active collector, if measurement is enabled. |
measurement_session(collector=None)
¶
Activate a collector for code running in this context.
Source code in src/anonymizer/measurement/session.py
@contextmanager
def measurement_session(collector: MeasurementCollector | None = None) -> Iterator[MeasurementCollector]:
"""Activate a collector for code running in this context."""
active = collector or MeasurementCollector()
token = _ACTIVE_COLLECTOR.set(active)
try:
yield active
finally:
_ACTIVE_COLLECTOR.reset(token)
configured_measurement_session(config)
¶
Activate and persist a collector when a measurement config is provided.
Source code in src/anonymizer/measurement/session.py
@contextmanager
def configured_measurement_session(config: MeasurementConfig | None) -> Iterator[MeasurementCollector | None]:
"""Activate and persist a collector when a measurement config is provided."""
if config is None:
yield None
return
sink = _JsonlMeasurementSink(config.output_path) if config.streaming else None
dd_trace_sink = None
if config.dd_trace != "none":
if config.dd_trace_path is None:
raise ValueError("dd_trace_path is required when dd_trace is enabled")
dd_trace_sink = _JsonlMeasurementSink(config.dd_trace_path)
dd_task_trace_sink = _JsonlMeasurementSink(config.dd_task_trace_path) if config.dd_task_trace_path else None
collector = MeasurementCollector(
run_id=config.run_id,
record_hash_key=config.record_hash_key,
record_level=config.record_level,
run_tags=config.run_tags,
record_sink=sink,
keep_records=config.keep_records,
dd_trace_mode=config.dd_trace,
dd_trace_sink=dd_trace_sink,
dd_task_trace_sink=dd_task_trace_sink,
fail_on_write_error=config.fail_on_write_error,
)
with measurement_session(collector):
body_error: BaseException | None = None
try:
yield collector
except BaseException as exc:
body_error = exc
raise
finally:
if config.streaming:
_close_collector_safely(config=config, collector=collector, body_error=body_error)
else:
write_error: BaseException | None = None
try:
_write_collector_safely(config=config, collector=collector, body_error=body_error)
except BaseException as exc:
write_error = exc
raise
finally:
_close_collector_safely(
config=config,
collector=collector,
body_error=body_error or write_error,
)
current_collector()
¶
Return the active collector, if measurement is enabled.
Source code in src/anonymizer/measurement/session.py
def current_collector() -> MeasurementCollector | None:
"""Return the active collector, if measurement is enabled."""
return _ACTIVE_COLLECTOR.get()