Skip to content

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()