Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
fd7a63b
feat: add SDK telemetry into emissions tracker
davidberenstein1957 Aug 20, 2026
f0fffb3
Merge remote-tracking branch 'origin/master' into HEAD
davidberenstein1957 Sep 23, 2026
cd8ccf6
fix(telemetry): address review on opt-out telemetry
davidberenstein1957 Sep 23, 2026
11ec802
Merge remote-tracking branch 'origin/master' into pr-1200-merge
davidberenstein1957 Sep 27, 2026
79ba393
refactor(telemetry): drop extensive tier, API token and coordinates
davidberenstein1957 Sep 27, 2026
50d377d
refactor(telemetry): keep only populated schema fields and cap strings
davidberenstein1957 Sep 27, 2026
e5a4f83
feat(carbonserver): add alembic migration for the telemetry table
davidberenstein1957 Sep 27, 2026
c2e01c5
feat(telemetry): send once per process and rate limit the endpoint
davidberenstein1957 Sep 27, 2026
dd772f5
test(telemetry): cover every changed telemetry line
davidberenstein1957 Sep 27, 2026
faeb497
test(telemetry): add end-to-end script checking rows in Postgres
davidberenstein1957 Sep 27, 2026
f608203
fix(telemetry): migrate pre-existing telemetry table instead of assum…
davidberenstein1957 Sep 27, 2026
9423cb7
fix(telemetry): fail closed to disabled on an unparseable telemetry_l…
davidberenstein1957 Sep 27, 2026
f7acc2d
fix(telemetry): make send_at_stop resilient to thread start failures
davidberenstein1957 Sep 27, 2026
d9b0627
fix(telemetry): truncate the payload timestamp to the hour
davidberenstein1957 Sep 27, 2026
a9d9c13
fix(telemetry): make the API rate limiter thread-safe
davidberenstein1957 Sep 27, 2026
35e6d09
fix(scripts): hard-fail e2e_telemetry.sh on Postgres or API startup p…
davidberenstein1957 Sep 27, 2026
2025c20
fix(telemetry): do not report the Canada geolocation fallback as a lo…
davidberenstein1957 Sep 27, 2026
992992d
fix(scripts): run e2e Postgres on a configurable host port
davidberenstein1957 Sep 27, 2026
2abce30
Potential fix for pull request finding 'CodeQL / Clear-text logging o…
davidberenstein1957 Sep 27, 2026
1db720e
fix(telemetry): write opt-out level to the global config by default
davidberenstein1957 Sep 27, 2026
2d388f8
fix(telemetry): print the one-time notice to stderr so it isn't hidden
davidberenstein1957 Sep 27, 2026
54b328d
fix(telemetry): read cudnn version from sys.modules, never import torch
davidberenstein1957 Sep 27, 2026
1ad9810
docs(telemetry): fix precedence order (env overrides config file)
davidberenstein1957 Sep 27, 2026
8ecf2e2
Review fixes
Sep 27, 2026
bc9c808
fix(telemetry): never send telemetry from OfflineEmissionsTracker
davidberenstein1957 Sep 27, 2026
5934207
fix(telemetry): treat empty telemetry_level as invalid, not unset
davidberenstein1957 Sep 27, 2026
3d94a71
fix(telemetry): stop reusing dashboard api_endpoint for telemetry URL
davidberenstein1957 Sep 27, 2026
15bf9de
test(config): drop leftover CODECARBON_TELEMETRY_PROJECT_TOKEN refere…
davidberenstein1957 Sep 27, 2026
fe01571
fix(migration): make telemetry.timestamp timezone-aware
davidberenstein1957 Sep 27, 2026
f0b6b4d
Merge remote-tracking branch 'origin/feat/add-telemetry' into pr-1200…
davidberenstein1957 Sep 27, 2026
4401307
feat(telemetry): show the notice every run until a level is chosen
davidberenstein1957 Sep 28, 2026
9005900
fix(telemetry): show the notice once, and skip the level prompt with …
davidberenstein1957 Sep 28, 2026
c42ea37
docs(telemetry): keep telemetry rows for at most 3 years
davidberenstein1957 Sep 28, 2026
442e781
feat(carbonserver): delete telemetry rows older than 3 years with a d…
davidberenstein1957 Sep 28, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

import uuid

from sqlalchemy import JSON, Boolean, Column, DateTime, Float, Integer, String
from sqlalchemy import Column, DateTime, Float, Integer, String
from sqlalchemy.dialects.postgresql import UUID

from carbonserver.database.database import Base
Expand All @@ -12,7 +12,7 @@ class Telemetry(Base):
__tablename__ = "telemetry"

id = Column(UUID(as_uuid=True), primary_key=True, index=True, default=uuid.uuid4)
timestamp = Column(DateTime, nullable=False)
timestamp = Column(DateTime(timezone=True), nullable=False)
telemetry_level = Column(String, nullable=False)

os = Column(String, nullable=True)
Expand All @@ -21,8 +21,6 @@ class Telemetry(Base):
region = Column(String, nullable=True)
cloud_provider = Column(String, nullable=True)
cloud_region = Column(String, nullable=True)
longitude = Column(Float, nullable=True)
latitude = Column(Float, nullable=True)

cpu_count = Column(Integer, nullable=True)
cpu_physical_count = Column(Integer, nullable=True)
Expand All @@ -38,64 +36,10 @@ class Telemetry(Base):

python_version = Column(String, nullable=True)
python_implementation = Column(String, nullable=True)
python_executable_hash = Column(String, nullable=True)
python_env_type = Column(String, nullable=True)
codecarbon_version = Column(String, nullable=True)
codecarbon_install_method = Column(String, nullable=True)

total_emissions_kg = Column(Float, nullable=True)
emissions_rate_kg_per_sec = Column(Float, nullable=True)
energy_consumed_kwh = Column(Float, nullable=True)
cpu_energy_kwh = Column(Float, nullable=True)
gpu_energy_kwh = Column(Float, nullable=True)
ram_energy_kwh = Column(Float, nullable=True)
duration_seconds = Column(Float, nullable=True)
cpu_utilization_avg = Column(Float, nullable=True)
gpu_utilization_avg = Column(Float, nullable=True)
ram_utilization_avg = Column(Float, nullable=True)

tracking_mode = Column(String, nullable=True)
api_mode = Column(String, nullable=True)
output_methods = Column(JSON, nullable=True)
hardware_tracked = Column(JSON, nullable=True)
task_tracking_used = Column(Boolean, nullable=True)
decorator_vs_context = Column(String, nullable=True)
measure_power_interval_secs = Column(Float, nullable=True)

hardware_detection_success = Column(Boolean, nullable=True)
rapl_available = Column(Boolean, nullable=True)
gpu_detection_method = Column(String, nullable=True)
first_measurement_time_ms = Column(Float, nullable=True)
tracking_overhead_percent = Column(Float, nullable=True)
errors_encountered = Column(JSON, nullable=True)
warning_count = Column(Integer, nullable=True)

ide_used = Column(String, nullable=True)
notebook_environment = Column(String, nullable=True)
ci_environment = Column(String, nullable=True)
python_package_manager = Column(String, nullable=True)
framework_detected = Column(String, nullable=True)

has_torch = Column(Boolean, nullable=True)
torch_version = Column(String, nullable=True)
has_transformers = Column(Boolean, nullable=True)
transformers_version = Column(String, nullable=True)
has_diffusers = Column(Boolean, nullable=True)
diffusers_version = Column(String, nullable=True)
has_tensorflow = Column(Boolean, nullable=True)
tensorflow_version = Column(String, nullable=True)
has_keras = Column(Boolean, nullable=True)
keras_version = Column(String, nullable=True)
has_pytorch_lightning = Column(Boolean, nullable=True)
pytorch_lightning_version = Column(String, nullable=True)
has_fastai = Column(Boolean, nullable=True)
fastai_version = Column(String, nullable=True)
ml_framework_primary = Column(String, nullable=True)

container_runtime = Column(String, nullable=True)
in_container = Column(Boolean, nullable=True)
host_machine_hash = Column(String, nullable=True)

def __repr__(self):
return (
f'<Telemetry(id="{self.id}", '
Expand Down
40 changes: 39 additions & 1 deletion carbonserver/carbonserver/api/routers/telemetry.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,12 @@
"""API router for handling telemetry data in the CarbonServer API."""

import threading
import time
from collections import defaultdict, deque
from uuid import UUID

from dependency_injector.wiring import Provide, inject
from fastapi import APIRouter, Depends
from fastapi import APIRouter, Depends, HTTPException, Request
from starlette import status

from carbonserver.api.schemas_telemetry import TelemetryCreate
Expand All @@ -12,9 +15,37 @@

TELEMETRY_ROUTER_TAGS = ["Telemetry"]

#: Accepted requests per client IP per window. The SDK sends once per process.
RATE_LIMIT = 60
RATE_WINDOW_SECONDS = 60.0
#: Forget every IP once this many are tracked, so memory stays bounded.
MAX_TRACKED_IPS = 10_000

# ponytail: in-memory, per-process limit. With more than one API instance or
# worker, each keeps its own counts; move to a proxy or Redis limit then.
_recent_requests: dict[str, deque] = defaultdict(deque)
# FastAPI runs sync path operations in a threadpool, so concurrent requests
# can call _rate_limited at once; without this lock two requests can race
# past the len(hits) >= RATE_LIMIT check and both be admitted.
_recent_requests_lock = threading.Lock()

router = APIRouter()


def _rate_limited(host: str) -> bool:
now = time.monotonic()
with _recent_requests_lock:
if host not in _recent_requests and len(_recent_requests) >= MAX_TRACKED_IPS:
_recent_requests.clear()
hits = _recent_requests[host]
while hits and now - hits[0] > RATE_WINDOW_SECONDS:
hits.popleft()
if len(hits) >= RATE_LIMIT:
return True
hits.append(now)
return False


@router.post(
"/telemetry",
tags=TELEMETRY_ROUTER_TAGS,
Expand All @@ -24,8 +55,15 @@
@inject
def add_telemetry(
telemetry: TelemetryCreate,
request: Request,
telemetry_service: TelemetryService = Depends(
Provide[ServerContainer.telemetry_service]
),
) -> UUID:
host = request.client.host if request.client else "unknown"
if _rate_limited(host):
raise HTTPException(
status_code=status.HTTP_429_TOO_MANY_REQUESTS,
detail="Too many telemetry requests",
)
return telemetry_service.add_telemetry(telemetry)
143 changes: 21 additions & 122 deletions carbonserver/carbonserver/api/schemas_telemetry.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,15 +2,17 @@

from datetime import datetime
from enum import Enum
from typing import List, Optional
from typing import Annotated, Optional

from pydantic import BaseModel, ConfigDict, Field, model_validator

#: Every free-text field is capped so one request cannot store unbounded data.
Str = Annotated[str, Field(max_length=256)]


class TelemetryLevel(str, Enum):
disabled = "disabled"
minimal = "minimal"
extensive = "extensive"


class TelemetryBase(BaseModel):
Expand All @@ -35,141 +37,38 @@ class TelemetryBase(BaseModel):
timestamp: datetime
telemetry_level: TelemetryLevel

os: Optional[str] = None
country_name: Optional[str] = None
os: Optional[Str] = None
country_name: Optional[Str] = None
country_iso_code: Optional[str] = Field(default=None, min_length=2, max_length=3)
region: Optional[str] = None
cloud_provider: Optional[str] = None
cloud_region: Optional[str] = None
longitude: Optional[float] = Field(default=None, ge=-180, le=180)
latitude: Optional[float] = Field(default=None, ge=-90, le=90)
region: Optional[Str] = None
cloud_provider: Optional[Str] = None
cloud_region: Optional[Str] = None

cpu_count: Optional[int] = Field(default=None, ge=0)
cpu_physical_count: Optional[int] = Field(default=None, ge=0)
cpu_model: Optional[str] = None
cpu_architecture: Optional[str] = None
cpu_model: Optional[Str] = None
cpu_architecture: Optional[Str] = None
gpu_count: Optional[int] = Field(default=None, ge=0)
gpu_model: Optional[str] = None
gpu_driver_version: Optional[str] = None
gpu_model: Optional[Str] = None
gpu_driver_version: Optional[Str] = None
gpu_memory_total_gb: Optional[float] = Field(default=None, ge=0)
ram_total_size_gb: Optional[float] = Field(default=None, ge=0)
cuda_version: Optional[str] = None
cudnn_version: Optional[str] = None
cuda_version: Optional[Str] = None
cudnn_version: Optional[Str] = None

python_version: Optional[str] = None
python_implementation: Optional[str] = None
python_executable_hash: Optional[str] = Field(
default=None, min_length=64, max_length=64
)
python_env_type: Optional[str] = None
codecarbon_version: Optional[str] = None
codecarbon_install_method: Optional[str] = None

total_emissions_kg: Optional[float] = Field(default=None, ge=0)
emissions_rate_kg_per_sec: Optional[float] = Field(default=None, ge=0)
energy_consumed_kwh: Optional[float] = Field(default=None, ge=0)
cpu_energy_kwh: Optional[float] = Field(default=None, ge=0)
gpu_energy_kwh: Optional[float] = Field(default=None, ge=0)
ram_energy_kwh: Optional[float] = Field(default=None, ge=0)
duration_seconds: Optional[float] = Field(default=None, ge=0)
cpu_utilization_avg: Optional[float] = Field(default=None, ge=0, le=100)
gpu_utilization_avg: Optional[float] = Field(default=None, ge=0, le=100)
ram_utilization_avg: Optional[float] = Field(default=None, ge=0, le=100)

tracking_mode: Optional[str] = None
api_mode: Optional[str] = None
output_methods: Optional[List[str]] = None
hardware_tracked: Optional[List[str]] = None
task_tracking_used: Optional[bool] = None
decorator_vs_context: Optional[str] = None
measure_power_interval_secs: Optional[float] = Field(default=None, ge=0)

hardware_detection_success: Optional[bool] = None
rapl_available: Optional[bool] = None
gpu_detection_method: Optional[str] = None
first_measurement_time_ms: Optional[float] = Field(default=None, ge=0)
tracking_overhead_percent: Optional[float] = Field(default=None, ge=0)
errors_encountered: Optional[List[str]] = None
warning_count: Optional[int] = Field(default=None, ge=0)

ide_used: Optional[str] = None
notebook_environment: Optional[str] = None
ci_environment: Optional[str] = None
python_package_manager: Optional[str] = None
framework_detected: Optional[str] = None

has_torch: Optional[bool] = None
torch_version: Optional[str] = None
has_transformers: Optional[bool] = None
transformers_version: Optional[str] = None
has_diffusers: Optional[bool] = None
diffusers_version: Optional[str] = None
has_tensorflow: Optional[bool] = None
tensorflow_version: Optional[str] = None
has_keras: Optional[bool] = None
keras_version: Optional[str] = None
has_pytorch_lightning: Optional[bool] = None
pytorch_lightning_version: Optional[str] = None
has_fastai: Optional[bool] = None
fastai_version: Optional[str] = None
ml_framework_primary: Optional[str] = None

container_runtime: Optional[str] = None
in_container: Optional[bool] = None
host_machine_hash: Optional[str] = None
python_version: Optional[Str] = None
python_implementation: Optional[Str] = None
python_env_type: Optional[Str] = None
codecarbon_version: Optional[Str] = None
codecarbon_install_method: Optional[Str] = None

@model_validator(mode="after")
def validate_telemetry_level(self):
def reject_disabled_level(self):
if self.telemetry_level == TelemetryLevel.disabled:
raise ValueError("Disabled telemetry must not be submitted")

if self.telemetry_level == TelemetryLevel.minimal:
extensive_fields = set(type(self).model_fields) - MINIMAL_TELEMETRY_FIELDS
submitted_extensive_fields = [
field
for field in extensive_fields
if getattr(self, field) not in (None, [], {})
]
if submitted_extensive_fields:
fields = ", ".join(sorted(submitted_extensive_fields))
raise ValueError(
f"Minimal telemetry cannot include extensive fields: {fields}"
)

return self


MINIMAL_TELEMETRY_FIELDS = {
"timestamp",
"telemetry_level",
"os",
"country_name",
"country_iso_code",
"region",
"cloud_provider",
"cloud_region",
"longitude",
"latitude",
"cpu_count",
"cpu_physical_count",
"cpu_model",
"cpu_architecture",
"gpu_count",
"gpu_model",
"gpu_driver_version",
"gpu_memory_total_gb",
"ram_total_size_gb",
"cuda_version",
"cudnn_version",
"python_version",
"python_implementation",
"python_executable_hash",
"python_env_type",
"codecarbon_version",
"codecarbon_install_method",
}


class TelemetryCreate(TelemetryBase):
pass

Expand Down
Loading
Loading