Skip to content

Commit 8689a8b

Browse files
czi-fsisendaadnanrhussain
authored andcommitted
chore: address PR comments
1 parent 754b809 commit 8689a8b

6 files changed

Lines changed: 118 additions & 50 deletions

File tree

sdks/python/src/learning_commons_evaluators/schemas/config.py

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,9 @@
1616

1717
DEFAULT_TELEMETRY_EVENTS_ENDPOINT = "https://api.learningcommons.org/evaluators-telemetry/v1/events"
1818

19+
# Shared per process so multiple :class:`EvaluatorConfig` instances derive the same client id.
20+
_PROCESS_CLIENT_ID_SEED = uuid.uuid4()
21+
1922
# --- LLM provider configs (for LLM calls in prompt steps) ---
2023

2124

@@ -112,8 +115,8 @@ class EvaluatorConfig:
112115

113116
# Temporary until we finalize the telemetry API key/client id strategy.
114117
#: UUID v5 namespace for deriving ``X-Client-ID`` when ``telemetry_partner_id`` is an API key.
115-
#: A new value is generated for each :class:`EvaluatorConfig` unless you pass one explicitly.
116-
client_id_seed: uuid.UUID = field(default_factory=uuid.uuid4)
118+
#: Defaults to a single per-process seed so all configs in one run share the same derived id.
119+
client_id_seed: uuid.UUID = field(default=_PROCESS_CLIENT_ID_SEED)
117120

118121

119122
def create_config(

sdks/python/src/learning_commons_evaluators/telemetry/__init__.py

Lines changed: 52 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -5,14 +5,16 @@
55
import asyncio
66
import threading
77
import uuid
8+
from concurrent.futures import ThreadPoolExecutor
9+
from datetime import datetime, timezone
810

911
import httpx
1012

1113
from learning_commons_evaluators.schemas.config import EvaluatorConfig
1214
from learning_commons_evaluators.schemas.evaluator import EvaluationInput
1315
from learning_commons_evaluators.schemas.metadata import EvaluationMetadata
1416
from learning_commons_evaluators.telemetry.adapter import evaluation_to_typescript_telemetry_event
15-
from learning_commons_evaluators.telemetry.utils import client_id_from_seed
17+
from learning_commons_evaluators.telemetry.utils import client_id_from_seed, iso_utc_z
1618

1719
__all__ = [
1820
"evaluation_to_typescript_telemetry_event",
@@ -21,6 +23,20 @@
2123
"should_send_telemetry",
2224
]
2325

26+
_TELEMETRY_EXECUTOR: ThreadPoolExecutor | None = None
27+
_TELEMETRY_EXECUTOR_LOCK = threading.Lock()
28+
29+
30+
def _get_telemetry_executor() -> ThreadPoolExecutor:
31+
global _TELEMETRY_EXECUTOR
32+
with _TELEMETRY_EXECUTOR_LOCK:
33+
if _TELEMETRY_EXECUTOR is None:
34+
_TELEMETRY_EXECUTOR = ThreadPoolExecutor(
35+
max_workers=2,
36+
thread_name_prefix="lc-telemetry",
37+
)
38+
return _TELEMETRY_EXECUTOR
39+
2440

2541
def should_send_telemetry(config: EvaluatorConfig) -> bool:
2642
"""Return True when telemetry is configured with a non-empty partner / client id."""
@@ -47,54 +63,59 @@ async def send_telemetry(
4763
if not should_send_telemetry(config):
4864
return
4965

50-
partner_id = config.telemetry.telemetry_partner_id
51-
assert (
52-
partner_id is not None
53-
) # for mypy: ``should_send_telemetry`` guarantees non-empty after strip.
54-
telemetry_partner_id = partner_id.strip()
55-
56-
event = evaluation_to_typescript_telemetry_event(evaluation_metadata, inp, config)
57-
payload = event.model_dump(mode="json", exclude_none=True)
58-
59-
api_key = telemetry_partner_id if not _is_uuid(telemetry_partner_id) else None
60-
client_id = (
61-
telemetry_partner_id
62-
if _is_uuid(telemetry_partner_id)
63-
else client_id_from_seed(telemetry_partner_id, config.client_id_seed)
64-
)
65-
66-
headers: dict[str, str] = {
67-
"Content-Type": "application/json",
68-
"X-Client-ID": client_id,
69-
}
70-
if api_key is not None:
71-
headers["X-API-Key"] = api_key
72-
7366
try:
67+
partner_id = config.telemetry.telemetry_partner_id
68+
assert (
69+
partner_id is not None
70+
) # for mypy: ``should_send_telemetry`` guarantees non-empty after strip.
71+
telemetry_partner_id = partner_id.strip()
72+
73+
event = evaluation_to_typescript_telemetry_event(evaluation_metadata, inp, config)
74+
# TS SDK sets timestamp at send time (`new Date().toISOString()`), not evaluation start.
75+
event = event.model_copy(update={"timestamp": iso_utc_z(datetime.now(timezone.utc))})
76+
payload = event.model_dump(mode="json", exclude_none=True)
77+
78+
api_key = telemetry_partner_id if not _is_uuid(telemetry_partner_id) else None
79+
client_id = (
80+
telemetry_partner_id
81+
if _is_uuid(telemetry_partner_id)
82+
else client_id_from_seed(telemetry_partner_id, config.client_id_seed)
83+
)
84+
85+
headers: dict[str, str] = {
86+
"Content-Type": "application/json",
87+
"X-Client-ID": client_id,
88+
}
89+
if api_key is not None:
90+
headers["X-API-Key"] = api_key
91+
7492
timeout = httpx.Timeout(5.0)
7593
async with httpx.AsyncClient(timeout=timeout) as client:
7694
response = await client.post(config.telemetry.endpoint, json=payload, headers=headers)
7795
if response.is_error:
96+
# Log status only; response bodies may echo input text or other sensitive data.
7897
config.logger.warning(
79-
"telemetry send failed: %s %s",
98+
"telemetry send failed: HTTP %s",
8099
response.status_code,
81-
response.text[:500],
82100
)
83-
except httpx.RequestError as e:
84-
config.logger.warning("telemetry send failed: %s", e)
101+
except Exception as e:
102+
# Log exception type only; ``str(e)`` may include payload fields (e.g. input_text).
103+
config.logger.warning(
104+
"telemetry send failed: %s",
105+
type(e).__qualname__,
106+
)
85107

86108

87109
def schedule_send_telemetry(
88110
evaluation_metadata: EvaluationMetadata,
89111
inp: EvaluationInput | None,
90112
config: EvaluatorConfig,
91113
) -> None:
92-
"""Fire-and-forget: run :func:`send_telemetry` on a daemon thread when telemetry is enabled."""
114+
"""Fire-and-forget: run :func:`send_telemetry` on a shared worker when telemetry is enabled."""
93115
if not should_send_telemetry(config):
94116
return
95117

96118
def _run() -> None:
97119
asyncio.run(send_telemetry(evaluation_metadata, inp, config))
98120

99-
thread = threading.Thread(target=_run, daemon=True)
100-
thread.start()
121+
_get_telemetry_executor().submit(_run)

sdks/python/src/learning_commons_evaluators/telemetry/adapter.py

Lines changed: 2 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,6 @@
1919

2020
from __future__ import annotations
2121

22-
from datetime import datetime, timezone
2322
from typing import Any
2423

2524
from learning_commons_evaluators.schemas.config import EvaluatorConfig, LLMProvider
@@ -41,11 +40,7 @@
4140
TelemetryStageDetail,
4241
TelemetryTokenUsage,
4342
)
44-
45-
46-
def _iso_utc_z(dt: datetime) -> str:
47-
s = dt.astimezone(timezone.utc).isoformat()
48-
return s.replace("+00:00", "Z")
43+
from learning_commons_evaluators.telemetry.utils import iso_utc_z
4944

5045

5146
def _map_status(status: Status) -> EvaluationTelemetryStatus:
@@ -179,7 +174,7 @@ def evaluation_to_typescript_telemetry_event(
179174
)
180175

181176
return TelemetryEvent(
182-
timestamp=_iso_utc_z(evaluation_metadata.timestamp),
177+
timestamp=iso_utc_z(evaluation_metadata.timestamp),
183178
sdk_version=evaluation_metadata.evaluator_metadata.sdk_version,
184179
evaluator_type=evaluation_metadata.evaluator_metadata.id,
185180
grade=grade,

sdks/python/src/learning_commons_evaluators/telemetry/utils.py

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,13 @@
33
from __future__ import annotations
44

55
import uuid
6+
from datetime import datetime, timezone
7+
8+
9+
def iso_utc_z(dt: datetime) -> str:
10+
"""Format *dt* as an ISO 8601 UTC string with a ``Z`` suffix (TS wire format)."""
11+
s = dt.astimezone(timezone.utc).isoformat()
12+
return s.replace("+00:00", "Z")
613

714

815
def client_id_from_seed(learning_commons_api_key: str, client_id_seed: uuid.UUID) -> str:

sdks/python/tests/schemas/test_config.py

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,11 @@ def test_create_config_sets_telemetry_partner_id(self):
7272
assert config.telemetry.endpoint == DEFAULT_TELEMETRY_EVENTS_ENDPOINT
7373
assert config.telemetry.send_full_input_with_telemetry is False
7474

75+
def test_default_client_id_seed_shared_across_configs(self):
76+
a = create_config(telemetry_partner_id="tid-a")
77+
b = create_config(telemetry_partner_id="tid-b")
78+
assert a.client_id_seed == b.client_id_seed
79+
7580
def test_create_config_telemetry_with_full_input_sets_flag(self):
7681
config = create_config_telemetry_with_full_input(telemetry_partner_id="tid")
7782
assert config.telemetry.telemetry_partner_id == "tid"

sdks/python/tests/telemetry/test_telemetry.py

Lines changed: 47 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
from __future__ import annotations
44

55
import asyncio
6+
from datetime import datetime, timezone
67
from unittest.mock import AsyncMock, MagicMock, patch
78

89
from learning_commons_evaluators import (
@@ -49,22 +50,21 @@ def test_true_with_partner_id(self):
4950

5051

5152
class TestScheduleSendTelemetry:
52-
def test_does_not_spawn_thread_when_partner_missing(self, config, evaluation_metadata):
53-
with patch("learning_commons_evaluators.telemetry.threading.Thread") as mock_thread:
53+
def test_does_not_submit_when_partner_missing(self, config, evaluation_metadata):
54+
with patch("learning_commons_evaluators.telemetry._get_telemetry_executor") as mock_exec:
5455
schedule_send_telemetry(evaluation_metadata, None, config)
55-
mock_thread.assert_not_called()
56+
mock_exec.assert_not_called()
5657

57-
@patch("learning_commons_evaluators.telemetry.threading.Thread")
58-
def test_starts_daemon_thread_when_no_running_loop(self, mock_thread, evaluation_metadata):
58+
@patch("learning_commons_evaluators.telemetry._get_telemetry_executor")
59+
def test_submits_to_shared_executor(self, mock_get_executor, evaluation_metadata):
5960
cfg = create_config(telemetry_partner_id="tid")
60-
mock_instance = MagicMock()
61-
mock_thread.return_value = mock_instance
61+
mock_executor = MagicMock()
62+
mock_get_executor.return_value = mock_executor
6263

6364
schedule_send_telemetry(evaluation_metadata, None, cfg)
6465

65-
mock_thread.assert_called_once()
66-
assert mock_thread.call_args.kwargs.get("daemon") is True
67-
mock_instance.start.assert_called_once()
66+
mock_get_executor.assert_called_once()
67+
mock_executor.submit.assert_called_once()
6868

6969

7070
class TestSendTelemetryHttp:
@@ -168,3 +168,40 @@ def test_send_telemetry_no_op_without_partner(
168168
):
169169
asyncio.run(send_telemetry(evaluation_metadata, None, config))
170170
mock_client_class.assert_not_called()
171+
172+
@patch("learning_commons_evaluators.telemetry.httpx.AsyncClient")
173+
def test_timestamp_is_set_at_send_time_not_evaluation_start(
174+
self, mock_client_class, evaluation_metadata
175+
):
176+
_mock_async_httpx_client(mock_client_class)
177+
cfg = create_config(telemetry_partner_id="tid")
178+
evaluation_metadata.timestamp = datetime(2020, 1, 1, tzinfo=timezone.utc)
179+
evaluation_metadata.status = Status.succeeded
180+
181+
before = datetime.now(timezone.utc)
182+
asyncio.run(send_telemetry(evaluation_metadata, None, cfg))
183+
after = datetime.now(timezone.utc)
184+
185+
payload = mock_client_class.return_value.__aenter__.return_value.post.call_args[1]["json"]
186+
sent = datetime.fromisoformat(payload["timestamp"].replace("Z", "+00:00"))
187+
assert before <= sent <= after
188+
189+
@patch(
190+
"learning_commons_evaluators.telemetry.evaluation_to_typescript_telemetry_event",
191+
side_effect=RuntimeError("adapter blew up"),
192+
)
193+
def test_send_telemetry_swallows_non_http_errors(self, _mock_adapter, evaluation_metadata):
194+
mock_logger = MagicMock()
195+
cfg = create_config(telemetry_partner_id="tid", logger=mock_logger)
196+
asyncio.run(send_telemetry(evaluation_metadata, None, cfg))
197+
mock_logger.warning.assert_called_once()
198+
199+
200+
class TestClientIdSeed:
201+
def test_same_partner_id_across_configs_yields_same_client_id(self):
202+
cfg_a = create_config(telemetry_partner_id="my-key")
203+
cfg_b = create_config(telemetry_partner_id="my-key")
204+
assert cfg_a.client_id_seed == cfg_b.client_id_seed
205+
assert client_id_from_seed("my-key", cfg_a.client_id_seed) == client_id_from_seed(
206+
"my-key", cfg_b.client_id_seed
207+
)

0 commit comments

Comments
 (0)