From e293e0cd0aec4728b9df1406eab15381d90962f2 Mon Sep 17 00:00:00 2001 From: David Berenstein Date: Wed, 19 Aug 2026 17:15:29 +0200 Subject: [PATCH 1/4] feat(fastapi): attribute per-request energy via fair-share windows Adds an ASGI middleware that gives each HTTP request its share of a long-running tracker's energy, plus the attribution model behind it. One tracker runs for the app's lifetime. Each completed sampling window (t_prev, t_now, dE) is split across the requests in flight during it, weighted by their overlap with the window and normalised by the sum of the weights. Windows with nothing in flight are recorded as unattributed. The invariant attributed + unattributed == settled holds exactly after every window, and is what the concurrency test pins down. Why not per-request start/stop energy snapshots: with N requests in flight each request observes the whole machine's delta, so the sum overcounts by roughly N - measured up to 88x at 100 concurrent requests. Fair-share weighting is the only split that conserves the run total. A request's share is only known one or more sampling windows after its response was sent, so results are reported then, via a callback. A request that never covered a completed window reports energy_kwh=None rather than zero: there is no honest number for it. Tracker side: add_energy_window_observer / remove_energy_window_observer expose the sampling windows, and http_request_emissions() scales the run's EmissionsData down to one attributed share using the run's accumulated component ratios and carbon intensity. Depends on #1374 (duration int -> float in the emissions schemas, and dropping the duration < 1 send guard) and #1375 (scheduler pause handling around tasks). Both are carried by their own PRs rather than duplicated here, so this should merge after them. Deliberately left out, to keep the diff reviewable: hardware-tier gating of which backends can resolve a sampling window, include/exclude path filtering (endpoint labelling is two lines inline), idle-baseline subtraction, per-endpoint aggregation, routing per-request rows into the tracker's own CSV/API output handlers, a lifespan helper, and a dedicated docs page. Each is additive on top of this and can follow if there is demand. Co-Authored-By: Claude Opus 5 (1M context) --- codecarbon/emissions_tracker.py | 62 ++++++ codecarbon/integrations/__init__.py | 1 + codecarbon/integrations/fastapi/__init__.py | 14 ++ .../integrations/fastapi/attribution.py | 199 ++++++++++++++++++ codecarbon/integrations/fastapi/middleware.py | 118 +++++++++++ docs/how-to/examples.md | 36 ++++ pyproject.toml | 2 + tests/integrations/test_fastapi.py | 147 +++++++++++++ uv.lock | 84 ++++++++ 9 files changed, 663 insertions(+) create mode 100644 codecarbon/integrations/__init__.py create mode 100644 codecarbon/integrations/fastapi/__init__.py create mode 100644 codecarbon/integrations/fastapi/attribution.py create mode 100644 codecarbon/integrations/fastapi/middleware.py create mode 100644 tests/integrations/test_fastapi.py diff --git a/codecarbon/emissions_tracker.py b/codecarbon/emissions_tracker.py index 96ed00c91..96921550e 100644 --- a/codecarbon/emissions_tracker.py +++ b/codecarbon/emissions_tracker.py @@ -296,6 +296,7 @@ def _initialize_runtime_state(self) -> None: self._tasks: Dict[str, Task] = {} self._active_task: Optional[str] = None self._active_task_emissions_at_start: Optional[EmissionsData] = None + self._window_observers: List[Callable[[float], None]] = [] self._hardware = [] self._hardware_initialized = False @@ -975,6 +976,66 @@ def _update_emissions(self) -> None: self._total_emissions += delta_emissions self._last_energy_covered = self._total_energy + def add_energy_window_observer(self, callback: Callable[[float], None]) -> None: + """Call ``callback(total_energy_kwh)`` after every completed sampling window. + + The callback runs on whichever thread took the sample (normally the + scheduler thread), so it must be cheap. Used by the FastAPI per-request + energy attribution to split each window's energy across the requests + that were in flight during it. + + Args: + callback: Receives the tracker's cumulative energy in kWh. + """ + self._window_observers.append(callback) + + def remove_energy_window_observer(self, callback: Callable[[float], None]) -> None: + """Remove a callback registered with :meth:`add_energy_window_observer`.""" + if callback in self._window_observers: + self._window_observers.remove(callback) + + def _notify_energy_window_observers(self) -> None: + for callback in self._window_observers: + try: + callback(self._total_energy.kWh) + except Exception: + logger.exception("CodeCarbon energy window observer failed") + + def http_request_emissions( + self, energy_kwh: float, duration_s: float + ) -> EmissionsData: + """Scale the run's emissions data down to one attributed energy share. + + The share is split into cpu/gpu/ram in the same proportions the run has + accumulated so far, and converted with the run's effective carbon + intensity, so per-request numbers stay consistent with the run total. + Over a single request the per-component split is not separately known, + and the run ratio is the best estimate available. + + Args: + energy_kwh: Energy attributed to the request, from + :class:`~codecarbon.integrations.fastapi.EnergyAttributor`. + duration_s: Wall-clock duration of the request. + """ + snapshot = self._prepare_emissions_data() + total = snapshot.energy_consumed or 0.0 + + def _share(component: float) -> float: + return energy_kwh * (component / total) if total else 0.0 + + emissions = energy_kwh * (self._total_emissions / total if total else 0.0) + return dataclasses.replace( + snapshot, + duration=duration_s, + emissions=emissions, + emissions_rate=emissions / duration_s if duration_s else 0.0, + cpu_energy=_share(snapshot.cpu_energy), + gpu_energy=_share(snapshot.gpu_energy), + ram_energy=_share(snapshot.ram_energy), + energy_consumed=energy_kwh, + water_consumed=_share(snapshot.water_consumed), + ) + def _prepare_emissions_data(self) -> EmissionsData: """ Prepare the emissions data to be sent to the API or written to a file. @@ -1274,6 +1335,7 @@ def _measure_power_and_energy(self) -> None: self._do_measurements() self._last_measured_time = time.perf_counter() + self._notify_energy_window_observers() self._measure_occurrence += 1 # Special case: metrics and api calls are sent every `api_call_interval` measures if ( diff --git a/codecarbon/integrations/__init__.py b/codecarbon/integrations/__init__.py new file mode 100644 index 000000000..9c5777a96 --- /dev/null +++ b/codecarbon/integrations/__init__.py @@ -0,0 +1 @@ +"""Optional integrations for frameworks and platforms.""" diff --git a/codecarbon/integrations/fastapi/__init__.py b/codecarbon/integrations/fastapi/__init__.py new file mode 100644 index 000000000..8722d97ea --- /dev/null +++ b/codecarbon/integrations/fastapi/__init__.py @@ -0,0 +1,14 @@ +"""FastAPI integration: per-request energy attribution middleware.""" + +from codecarbon.integrations.fastapi.attribution import EnergyAttributor, RequestEnergy +from codecarbon.integrations.fastapi.middleware import ( + CodeCarbonMiddleware, + add_codecarbon_middleware, +) + +__all__ = [ + "CodeCarbonMiddleware", + "EnergyAttributor", + "RequestEnergy", + "add_codecarbon_middleware", +] diff --git a/codecarbon/integrations/fastapi/attribution.py b/codecarbon/integrations/fastapi/attribution.py new file mode 100644 index 000000000..eb17d30bc --- /dev/null +++ b/codecarbon/integrations/fastapi/attribution.py @@ -0,0 +1,199 @@ +"""Fair-share per-request energy attribution. + +Each completed sampling window ``(t_prev, t_now, dE)`` is split across the +requests that were in flight during it, weighted by their overlap with the +window and normalised **by the sum of the weights**. Windows with nothing in +flight go entirely to ``unattributed_kwh``. The invariant is:: + + attributed_kwh + unattributed_kwh == settled_kwh + +exactly, after every window. That is the property the tests pin down. + +Start/stop energy snapshots per request cannot do this: with N requests in +flight each one sees the whole machine's delta, so the sum overcounts by +roughly N (measured up to 88x at 100 concurrent requests). +""" + +from __future__ import annotations + +import threading +import time +from collections.abc import Callable +from dataclasses import dataclass +from typing import Any + +from codecarbon.external.logger import logger + + +@dataclass(frozen=True) +class RequestEnergy: + """One request's finished attribution. + + ``energy_kwh`` is ``None`` when the request never covered a completed + sampling window: there is no honest number, and zero would be a lie. + """ + + endpoint: str + energy_kwh: float | None + duration_s: float + #: Completed sampling windows this request overlapped. + windows: int + #: Mean number of requests it competed against, window-weighted. + mean_concurrency: float | None + + +@dataclass +class _InFlight: + """Mutable per-request state, freed as soon as the request emits.""" + + endpoint: str + start: float + end: float | None = None + energy: float = 0.0 + windows: int = 0 + concurrency_sum: float = 0.0 + #: Called with the :class:`RequestEnergy` when this request resolves. + on_resolved: Callable[[RequestEnergy], None] | None = None + + +class EnergyAttributor: + """Splits each sampling window's energy across the requests in flight.""" + + def __init__(self) -> None: + self._in_flight: dict[int, _InFlight] = {} + # begin/end run on the event-loop thread, on_window on the tracker's + # scheduler thread. Never held across an on_resolved callback. + self._lock = threading.Lock() + #: Running sum of everything handed to requests, kWh. + self.attributed_kwh = 0.0 + #: Energy from windows with nothing in flight, kWh. + self.unattributed_kwh = 0.0 + #: Energy taken in from closed windows. ``attributed + unattributed == + #: settled`` holds exactly after every window; it is below the tracker's + #: run total by whatever a wrapped counter dropped + #: (``windows_skipped``) plus the final unsampled partial window. + self.settled_kwh = 0.0 + self.windows_settled = 0 + #: Windows where the energy counter went backwards (RAPL wrap/reset). + self.windows_skipped = 0 + self._t_prev = time.perf_counter() + self._e_prev = 0.0 + + def reset_window(self, total_energy_kwh: float = 0.0) -> None: + """Anchor the first window at now. Call when the tracker starts.""" + self._t_prev = time.perf_counter() + self._e_prev = total_energy_kwh + + def begin(self, endpoint: str) -> _InFlight: + """Start weighting a request. Returns the handle to pass to :meth:`end`.""" + state = _InFlight(endpoint=endpoint, start=time.perf_counter()) + with self._lock: + self._in_flight[id(state)] = state + return state + + def end(self, state: _InFlight) -> None: + """Stamp the request finished. + + Deliberately does **not** settle. The request stays weighted until the + next real sample closes, because at response time the machine's power + over the last partial window is genuinely unknown - settling here would + drop that energy into a zero-width window and silently lose it. + """ + with self._lock: + state.end = time.perf_counter() + + def close(self) -> None: + """Emit every in-flight request as-is. Call after the tracker stops.""" + with self._lock: + pending = list(self._in_flight.values()) + self._in_flight.clear() + for state in pending: + self._emit(state) + + def on_window(self, total_energy_kwh: float) -> None: + """Close a sampling window with the tracker's cumulative energy. + + Wired to + :meth:`~codecarbon.emissions_tracker.BaseEmissionsTracker.add_energy_window_observer`, + so it is only ever called from a real hardware sample. + """ + with self._lock: + self._settle(total_energy_kwh) + finished = [ + self._in_flight.pop(key) + for key, state in list(self._in_flight.items()) + if state.end is not None + ] + # Emitted outside the lock: on_resolved is user code and may call begin. + for state in finished: + self._emit(state) + + def _settle(self, total_energy_kwh: float) -> None: + """Split one window. Caller must hold ``self._lock``.""" + now = time.perf_counter() + w0, w1 = self._t_prev, now + width = w1 - w0 + delta = total_energy_kwh - self._e_prev + if width <= 0: + self._t_prev, self._e_prev = now, total_energy_kwh + return + if delta < 0: + # Counter wraparound or reset: no honest way to split a negative. + self.windows_skipped += 1 + self._t_prev, self._e_prev = now, total_energy_kwh + return + + states: list[_InFlight] = [] + weights: list[float] = [] + for state in self._in_flight.values(): + lo = max(state.start, w0) + hi = min(state.end if state.end is not None else w1, w1) + if hi - lo <= 0: + continue + weights.append(hi - lo) + states.append(state) + + if not states: + self.unattributed_kwh += delta + else: + total_weight = sum(weights) + for state, weight in zip(states, weights): + share = delta * (weight / total_weight) + state.energy += share + state.windows += 1 + state.concurrency_sum += len(states) + self.attributed_kwh += share + + # Banked only once the split succeeded. The caller swallows exceptions, + # so advancing the cursor first would drop this window's energy from + # settled_kwh and break attributed + unattributed == settled. + self.windows_settled += 1 + self.settled_kwh += delta + self._t_prev, self._e_prev = now, total_energy_kwh + + def _emit(self, state: _InFlight) -> None: + result = RequestEnergy( + endpoint=state.endpoint, + energy_kwh=state.energy if state.windows else None, + duration_s=(state.end or time.perf_counter()) - state.start, + windows=state.windows, + mean_concurrency=( + state.concurrency_sum / state.windows if state.windows else None + ), + ) + if state.on_resolved is not None: + try: + state.on_resolved(result) + except Exception: + logger.exception("CodeCarbon attribution callback failed") + + def report(self) -> dict[str, Any]: + """Run-level accounting, for checking what the split did.""" + return { + "attributed_kwh": self.attributed_kwh, + "unattributed_kwh": self.unattributed_kwh, + "settled_kwh": self.settled_kwh, + "windows_settled": self.windows_settled, + "windows_skipped": self.windows_skipped, + "in_flight": len(self._in_flight), # racy read, reporting only + } diff --git a/codecarbon/integrations/fastapi/middleware.py b/codecarbon/integrations/fastapi/middleware.py new file mode 100644 index 000000000..cc82476e5 --- /dev/null +++ b/codecarbon/integrations/fastapi/middleware.py @@ -0,0 +1,118 @@ +"""FastAPI/Starlette middleware for per-request emissions attribution.""" + +from __future__ import annotations + +import functools +from collections.abc import Callable + +from starlette.types import ASGIApp, Message, Receive, Scope, Send + +from codecarbon.emissions_tracker import BaseEmissionsTracker +from codecarbon.external.logger import logger +from codecarbon.integrations.fastapi.attribution import EnergyAttributor, RequestEnergy +from codecarbon.output_methods.emissions_data import EmissionsData + + +def log_request( + energy: RequestEnergy, emissions: EmissionsData | None, status_code: int +) -> None: + """Default ``on_request`` handler; logs via the ``codecarbon`` logger.""" + logger.info( + "CodeCarbon %s: energy=%s kWh emissions=%s kg CO2 status=%s", + energy.endpoint, + energy.energy_kwh, + getattr(emissions, "emissions", None), + status_code, + ) + + +class CodeCarbonMiddleware: + """Attributes a running tracker's energy to each HTTP request. + + A request's number is only known one or more sampling windows *after* its + response was sent, so ``on_request`` is called then, from the tracker's + scheduler thread. Keep it cheap and non-blocking. + + Args: + app: Inner ASGI application. + tracker: A tracker already started with ``tracker.start()``. + on_request: Callback ``(RequestEnergy, EmissionsData | None, status_code)``. + ``None`` disables reporting. + """ + + def __init__( + self, + app: ASGIApp, + *, + tracker: BaseEmissionsTracker, + on_request: ( + Callable[[RequestEnergy, EmissionsData | None, int], None] | None + ) = log_request, + ) -> None: + self.app = app + self.tracker = tracker + self.on_request = on_request + self.attributor = EnergyAttributor() + self.attributor.reset_window(tracker._total_energy.kWh) + tracker.add_energy_window_observer(self.attributor.on_window) + + def close(self) -> None: + """Stop attributing and emit whatever is still in flight.""" + self.tracker.remove_energy_window_observer(self.attributor.on_window) + self.attributor.close() + + async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None: + """ASGI entrypoint.""" + if scope["type"] != "http": + await self.app(scope, receive, send) + return + + status_code = 500 + + async def send_wrapper(message: Message) -> None: + nonlocal status_code + if message["type"] == "http.response.start": + status_code = message["status"] + await send(message) + + state = self.attributor.begin(_endpoint(scope)) + try: + await self.app(scope, receive, send_wrapper) + finally: + # The route template only lands in the scope once Starlette's + # router has run, so the endpoint can only be named here. + state.endpoint = _endpoint(scope) + state.on_resolved = functools.partial(self._resolved, status_code) + self.attributor.end(state) + + def _resolved(self, status_code: int, energy: RequestEnergy) -> None: + if self.on_request is None: + return + emissions = ( + self.tracker.http_request_emissions(energy.energy_kwh, energy.duration_s) + if energy.energy_kwh is not None + else None + ) + self.on_request(energy, emissions, status_code) + + +def _endpoint(scope: Scope) -> str: + route = scope.get("route") + path = getattr(route, "path", None) or scope.get("path", "") + return f"{scope.get('method', '')} {path}".strip() + + +def add_codecarbon_middleware(app, **kwargs) -> None: + """Register :class:`CodeCarbonMiddleware` and expose it on ``app.state``. + + Starlette builds the middleware stack on startup, so the instance that + actually serves requests is only knowable from inside its constructor. + ``app.state.codecarbon_middleware.close()`` on shutdown. + """ + + class _Registered(CodeCarbonMiddleware): + def __init__(self, asgi_app: ASGIApp, **kw) -> None: + super().__init__(asgi_app, **kw) + app.state.codecarbon_middleware = self + + app.add_middleware(_Registered, **kwargs) diff --git a/docs/how-to/examples.md b/docs/how-to/examples.md index 06f084aca..2763cd2c9 100644 --- a/docs/how-to/examples.md +++ b/docs/how-to/examples.md @@ -157,3 +157,39 @@ tracker.start_task("training") tracker.stop_task() tracker.stop() ``` + +### Track FastAPI Requests + +One tracker runs for the app lifetime; the middleware splits each of its +sampling windows across the requests that were in flight during that window, +weighted by overlap. Per-request start/stop snapshots cannot be used here: +with N requests in flight each one would see the whole machine's delta, so +the sum overcounts by roughly N. + +```python +from contextlib import asynccontextmanager + +from fastapi import FastAPI + +from codecarbon import EmissionsTracker +from codecarbon.integrations.fastapi import add_codecarbon_middleware + + +@asynccontextmanager +async def lifespan(app: FastAPI): + tracker = EmissionsTracker(allow_multiple_runs=True) + tracker.start() + add_codecarbon_middleware(app, tracker=tracker) + yield + tracker.stop() + app.state.codecarbon_middleware.close() + + +app = FastAPI(lifespan=lifespan) +``` + +A request's share is only known one or more sampling windows *after* its +response was sent, so the `on_request(energy, emissions, status_code)` +callback fires then, on the tracker's scheduler thread — keep it cheap. +`energy.energy_kwh` is `None` when the request was shorter than the gap +between two samples: there is no honest number, and zero would be a lie. diff --git a/pyproject.toml b/pyproject.toml index 795be29b1..1205e8e39 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -95,6 +95,8 @@ dev = [ "jsonschema", # For BoAmps schema validation tests "mktestdocs", # For testing documentation code blocks "scikit-learn", # For documentation examples and tests + "fastapi>=0.100", # For testing the FastAPI middleware + "httpx", # For fastapi.testclient ] doc = [ "requests", diff --git a/tests/integrations/test_fastapi.py b/tests/integrations/test_fastapi.py new file mode 100644 index 000000000..aa7fb4e1a --- /dev/null +++ b/tests/integrations/test_fastapi.py @@ -0,0 +1,147 @@ +"""Tests for the FastAPI per-request energy attribution.""" + +import threading +import time + +import pytest +from fastapi import FastAPI +from fastapi.testclient import TestClient + +from codecarbon.emissions_tracker import EmissionsTracker +from codecarbon.integrations.fastapi import ( + EnergyAttributor, + RequestEnergy, + add_codecarbon_middleware, +) + + +def _invariant(attributor: EnergyAttributor) -> None: + report = attributor.report() + assert report["attributed_kwh"] + report["unattributed_kwh"] == pytest.approx( + report["settled_kwh"], rel=1e-12, abs=1e-15 + ) + + +def test_idle_windows_are_unattributed(): + attributor = EnergyAttributor() + attributor.reset_window(0.0) + attributor.on_window(1.0) + assert attributor.unattributed_kwh == 1.0 + assert attributor.attributed_kwh == 0.0 + _invariant(attributor) + + +def test_window_energy_splits_by_overlap(): + attributor = EnergyAttributor() + attributor.reset_window(0.0) + early = attributor.begin("GET /a") + time.sleep(0.02) + late = attributor.begin("GET /b") + time.sleep(0.02) + attributor.on_window(1.0) + + # `early` overlapped roughly twice as much of the window as `late`. + assert early.energy > late.energy + assert early.energy + late.energy == pytest.approx(1.0) + _invariant(attributor) + + +def test_backwards_counter_is_skipped_not_split(): + attributor = EnergyAttributor() + attributor.reset_window(5.0) + attributor.begin("GET /a") + attributor.on_window(1.0) # RAPL wrap + assert attributor.windows_skipped == 1 + assert attributor.attributed_kwh == 0.0 + _invariant(attributor) + + +def test_unresolved_request_reports_no_energy(): + """A request that never covered a window gets None, not zero.""" + attributor = EnergyAttributor() + attributor.reset_window(0.0) + results = [] + state = attributor.begin("GET /fast") + state.on_resolved = results.append + attributor.end(state) + attributor.close() + + assert [r.energy_kwh for r in results] == [None] + + +def test_invariant_holds_under_concurrency(): + """The core property: nothing is created or lost by the split.""" + attributor = EnergyAttributor() + attributor.reset_window(0.0) + results: list[RequestEnergy] = [] # list.append is atomic under the GIL + stop = threading.Event() + energy = 0.0 + + def sampler(): + nonlocal energy + while not stop.is_set(): + energy += 0.001 + attributor.on_window(energy) + _invariant(attributor) + time.sleep(0.002) + + def requester(i: int): + for _ in range(20): + state = attributor.begin(f"GET /{i % 3}") + state.on_resolved = results.append + time.sleep(0.001) + attributor.end(state) + + sampler_thread = threading.Thread(target=sampler) + sampler_thread.start() + workers = [threading.Thread(target=requester, args=(i,)) for i in range(8)] + for w in workers: + w.start() + for w in workers: + w.join() + stop.set() + sampler_thread.join() + attributor.close() + + _invariant(attributor) + assert len(results) == 8 * 20 + resolved = [r.energy_kwh for r in results if r.energy_kwh is not None] + assert resolved, "no request ever covered a window" + assert sum(resolved) == pytest.approx(attributor.attributed_kwh) + + +def test_end_to_end_through_a_real_tracker(): + tracker = EmissionsTracker( + measure_power_secs=0.5, output_methods=[], allow_multiple_runs=True + ) + tracker.start() + seen: list[tuple] = [] + app = FastAPI() + add_codecarbon_middleware( + app, + tracker=tracker, + on_request=lambda energy, emissions, status: seen.append( + (energy, emissions, status) + ), + ) + + @app.get("/work/{n}") + def work(n: int): + time.sleep(0.6) # long enough to span a sampling window + return {"n": n} + + try: + with TestClient(app) as client: + assert client.get("/work/1").status_code == 200 + time.sleep(1.2) # let a window close and resolve the request + finally: + tracker.stop() + app.state.codecarbon_middleware.close() + + assert seen, "no request was reported" + energy, emissions, status = seen[0] + assert status == 200 + assert energy.endpoint == "GET /work/{n}" + assert energy.energy_kwh is not None and energy.energy_kwh > 0 + assert emissions.energy_consumed == energy.energy_kwh + assert emissions.duration == pytest.approx(energy.duration_s) diff --git a/uv.lock b/uv.lock index 7830e8098..b1c19ac9a 100644 --- a/uv.lock +++ b/uv.lock @@ -35,6 +35,20 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/99/91/8acff4f5e50511b911bbccb72b8628a49c68ce14148cd9f6431094859a90/annotated_types-0.8.0-py3-none-any.whl", hash = "sha256:f072f4d804ea359e4eaf198b1af7a8b0943881a87f31bb764f8bf219bb9419e0", size = 13427, upload-time = "2026-07-23T20:16:12.938Z" }, ] +[[package]] +name = "anyio" +version = "4.14.2" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "exceptiongroup", marker = "python_full_version < '3.11'" }, + { name = "idna" }, + { name = "typing-extensions", marker = "python_full_version < '3.13'" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/61/cc/a381afa6efea9f496eff839d4a6a1aed3bfafc7b3ab4b0d1b243a12573dd/anyio-4.14.2.tar.gz", hash = "sha256:cfa139f3ed1a23ee8f88a145ddb5ac7605b8bbfd8592baacd7ce3d8bb4313c7f", size = 260176, upload-time = "2026-07-12T20:29:07.082Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/da/35/f2287558c17e29fafc8ef3daf819bb9834061cfa43bff8014f7df7f63bdc/anyio-4.14.2-py3-none-any.whl", hash = "sha256:9f505dda5ac9f0c8309b5e8bd445a8c2bf7246f3ce950121e45ea15bc41d1494", size = 125813, upload-time = "2026-07-12T20:29:05.763Z" }, +] + [[package]] name = "arrow" version = "1.4.0" @@ -458,6 +472,8 @@ viz-legacy = [ dev = [ { name = "black" }, { name = "bumpver" }, + { name = "fastapi" }, + { name = "httpx" }, { name = "jsonschema" }, { name = "logfire" }, { name = "mktestdocs" }, @@ -515,6 +531,8 @@ provides-extras = ["carbonboard", "viz-legacy"] dev = [ { name = "black" }, { name = "bumpver" }, + { name = "fastapi", specifier = ">=0.100" }, + { name = "httpx" }, { name = "jsonschema" }, { name = "logfire", specifier = ">=1.0.1" }, { name = "mktestdocs" }, @@ -791,6 +809,22 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/c1/ea/53f2148663b321f21b5a606bd5f191517cf40b7072c0497d3c92c4a13b1e/executing-2.2.1-py2.py3-none-any.whl", hash = "sha256:760643d3452b4d777d295bb167ccc74c64a81df23fb5e08eff250c425a4b2017", size = 28317, upload-time = "2025-09-01T09:48:08.5Z" }, ] +[[package]] +name = "fastapi" +version = "0.141.1" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "annotated-doc" }, + { name = "pydantic" }, + { name = "starlette" }, + { name = "typing-extensions" }, + { name = "typing-inspection" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/8a/02/91e3416a8fdd715abb903a952a6bec7cdd8d14eed55d415fc8595524c319/fastapi-0.141.1.tar.gz", hash = "sha256:e8822fc40db1e1858054d7a949a888695bc9bdce70139178e33bd2871a453ca1", size = 425799, upload-time = "2026-07-29T17:18:05.568Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/cb/03/10388a42375ee7e4ac9b94eb2c5c569c8b5795e377e701c9ac3ad63de890/fastapi-0.141.1-py3-none-any.whl", hash = "sha256:bfb91aa2d334c61cb35ba9a116fc123b3d3df31640b801cf57a7a78ec3f603b3", size = 131954, upload-time = "2026-07-29T17:18:04.364Z" }, +] + [[package]] name = "filelock" version = "3.32.2" @@ -862,6 +896,43 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/e4/d3/5268aeabf2ad82658c4e2ff3a060648d0f02f3926cb53247c0e4d0dab49e/griffelib-2.1.0-py3-none-any.whl", hash = "sha256:cc7b3d2d2865ad0b909fcc38086e3f554b5ea7acbaa7bbb7ecaa3f5dfb7d9f00", size = 142560, upload-time = "2026-06-19T12:05:38.742Z" }, ] +[[package]] +name = "h11" +version = "0.16.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/01/ee/02a2c011bdab74c6fb3c75474d40b3052059d95df7e73351460c8588d963/h11-0.16.0.tar.gz", hash = "sha256:4e35b956cf45792e4caa5885e69fba00bdbc6ffafbfa020300e549b208ee5ff1", size = 101250, upload-time = "2025-04-24T03:35:25.427Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/04/4b/29cac41a4d98d144bf5f6d33995617b185d14b22401f75ca86f384e87ff1/h11-0.16.0-py3-none-any.whl", hash = "sha256:63cf8bbe7522de3bf65932fda1d9c2772064ffb3dae62d55932da54b31cb6c86", size = 37515, upload-time = "2025-04-24T03:35:24.344Z" }, +] + +[[package]] +name = "httpcore" +version = "1.0.9" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "certifi" }, + { name = "h11" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/06/94/82699a10bca87a5556c9c59b5963f2d039dbd239f25bc2a63907a05a14cb/httpcore-1.0.9.tar.gz", hash = "sha256:6e34463af53fd2ab5d807f399a9b45ea31c3dfa2276f15a2c3f00afff6e176e8", size = 85484, upload-time = "2025-04-24T22:06:22.219Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/7e/f5/f66802a942d491edb555dd61e3a9961140fd64c90bce1eafd741609d334d/httpcore-1.0.9-py3-none-any.whl", hash = "sha256:2d400746a40668fc9dec9810239072b40b4484b640a8c38fd654a024c7a1bf55", size = 78784, upload-time = "2025-04-24T22:06:20.566Z" }, +] + +[[package]] +name = "httpx" +version = "0.28.1" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "anyio" }, + { name = "certifi" }, + { name = "httpcore" }, + { name = "idna" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/b1/df/48c586a5fe32a0f01324ee087459e112ebb7224f646c0b5023f5e79e9956/httpx-0.28.1.tar.gz", hash = "sha256:75e98c5f16b0f35b567856f597f06ff2270a374470a5c2392242528e3e3e42fc", size = 141406, upload-time = "2024-12-06T15:37:23.222Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/2a/39/e50c7c3a983047577ee07d2a9e53faf5a69493943ec3f6a384bdc792deb2/httpx-0.28.1-py3-none-any.whl", hash = "sha256:d909fcccc110f8c7faf814ca82a9a4d816bc5a6dbfea25d6591d6985b8ba59ad", size = 73517, upload-time = "2024-12-06T15:37:21.509Z" }, +] + [[package]] name = "identify" version = "2.6.19" @@ -3192,6 +3263,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/0f/2c/437fe806897c2d6cfdc3ee43a18da8bf8e568530a4ae9bac781541ca9896/soupsieve-2.9.1-py3-none-any.whl", hash = "sha256:4f4477399246b7a0c720a88ca2454b11cd6bb9ae4c9d170140786e916776c14c", size = 37404, upload-time = "2026-07-21T16:57:16.421Z" }, ] +[[package]] +name = "starlette" +version = "1.6.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "anyio" }, + { name = "typing-extensions", marker = "python_full_version < '3.13'" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/b5/b4/205b0d5241d934e8add0c38aa924c4f9fb7330834ff11e5444db964ec3f9/starlette-1.6.0.tar.gz", hash = "sha256:d4e3ac5e546444960c710297a3c9fc3f7ebae1b7e963f3d36173b49da535be9b", size = 2716969, upload-time = "2026-08-08T18:27:57.512Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/c8/cb/6a6a47d5b464bd08695d254f3da6e7986cc70c9fa5d778eda57538edfe56/starlette-1.6.0-py3-none-any.whl", hash = "sha256:a86dd39d14bb45f85a3d18525215a9ef0cfd1f192ac793220e72598c90335f0c", size = 75969, upload-time = "2026-08-08T18:27:56.196Z" }, +] + [[package]] name = "taskipy" version = "1.14.1" From 6819b5042f51e4c407b9b94f3a87ffd519716e49 Mon Sep 17 00:00:00 2001 From: David Berenstein Date: Wed, 23 Sep 2026 16:50:11 +0900 Subject: [PATCH 2/4] fix(fastapi): address review on per-request attribution - Middleware is added at module level; the tracker is passed in or read from app.state.codecarbon_tracker, so the documented lifespan pattern no longer raises "Cannot add middleware after an application has started". Drop add_codecarbon_middleware and its per-call subclass. - Only record requests while the tracker runs; settle pending requests when it stops/changes and on lifespan shutdown (no unbounded growth). - Compute carbon intensity once per window and report emissions_kg as energy x intensity; remove http_request_emissions from the core tracker, so the scheduler thread no longer mutates run totals. - Iterate over a copy of the window observers. - Add a codecarbon[fastapi] extra with an install hint on ImportError. - Default log_request logs at DEBUG. - Tests: offline/fake trackers, no sleeps; cover the documented lifespan pattern, never-started tracker, 500 on raise, close on shutdown. Co-Authored-By: Claude Opus 5.5 (1M context) --- codecarbon/emissions_tracker.py | 46 ++----- codecarbon/integrations/fastapi/__init__.py | 6 +- codecarbon/integrations/fastapi/middleware.py | 104 ++++++++++----- docs/how-to/examples.md | 18 ++- pyproject.toml | 3 + tests/integrations/test_fastapi.py | 126 ++++++++++++++---- uv.lock | 46 ++++--- 7 files changed, 225 insertions(+), 124 deletions(-) diff --git a/codecarbon/emissions_tracker.py b/codecarbon/emissions_tracker.py index 8c94741eb..f14a5ffdd 100644 --- a/codecarbon/emissions_tracker.py +++ b/codecarbon/emissions_tracker.py @@ -1024,46 +1024,26 @@ def remove_energy_window_observer(self, callback: Callable[[float], None]) -> No self._window_observers.remove(callback) def _notify_energy_window_observers(self) -> None: - for callback in self._window_observers: + # Copy: an observer may be removed from another thread mid-iteration. + for callback in tuple(self._window_observers): try: callback(self._total_energy.kWh) except Exception: logger.exception("CodeCarbon energy window observer failed") - def http_request_emissions( - self, energy_kwh: float, duration_s: float - ) -> EmissionsData: - """Scale the run's emissions data down to one attributed energy share. + def _carbon_intensity_kg_per_kwh(self) -> float: + """Current carbon intensity, kg CO2eq per kWh, without touching run totals. - The share is split into cpu/gpu/ram in the same proportions the run has - accumulated so far, and converted with the run's effective carbon - intensity, so per-request numbers stay consistent with the run total. - Over a single request the per-component split is not separately known, - and the run ratio is the best estimate available. - - Args: - energy_kwh: Energy attributed to the request, from - :class:`~codecarbon.integrations.fastapi.EnergyAttributor`. - duration_s: Wall-clock duration of the request. + Used by the FastAPI integration to convert each request's energy share, + once per sampling window, off the tracker's accumulated state. """ - snapshot = self._prepare_emissions_data() - total = snapshot.energy_consumed or 0.0 - - def _share(component: float) -> float: - return energy_kwh * (component / total) if total else 0.0 - - emissions = energy_kwh * (self._total_emissions / total if total else 0.0) - return dataclasses.replace( - snapshot, - duration=duration_s, - emissions=emissions, - emissions_rate=emissions / duration_s if duration_s else 0.0, - cpu_energy=_share(snapshot.cpu_energy), - gpu_energy=_share(snapshot.gpu_energy), - ram_energy=_share(snapshot.ram_energy), - energy_consumed=energy_kwh, - water_consumed=_share(snapshot.water_consumed), - ) + self._ensure_geo_metadata() + self._ensure_emissions_engine() + one_kwh = Energy.from_energy(kWh=1) + cloud: CloudMetadata = self._get_cloud_metadata() + if cloud.is_on_private_infra: + return self._emissions.get_private_infra_emissions(one_kwh, self._geo) + return self._emissions.get_cloud_emissions(one_kwh, cloud, self._geo) def _prepare_emissions_data(self) -> EmissionsData: """ diff --git a/codecarbon/integrations/fastapi/__init__.py b/codecarbon/integrations/fastapi/__init__.py index 8722d97ea..dd1969761 100644 --- a/codecarbon/integrations/fastapi/__init__.py +++ b/codecarbon/integrations/fastapi/__init__.py @@ -1,14 +1,10 @@ """FastAPI integration: per-request energy attribution middleware.""" from codecarbon.integrations.fastapi.attribution import EnergyAttributor, RequestEnergy -from codecarbon.integrations.fastapi.middleware import ( - CodeCarbonMiddleware, - add_codecarbon_middleware, -) +from codecarbon.integrations.fastapi.middleware import CodeCarbonMiddleware __all__ = [ "CodeCarbonMiddleware", "EnergyAttributor", "RequestEnergy", - "add_codecarbon_middleware", ] diff --git a/codecarbon/integrations/fastapi/middleware.py b/codecarbon/integrations/fastapi/middleware.py index cc82476e5..e44ccc527 100644 --- a/codecarbon/integrations/fastapi/middleware.py +++ b/codecarbon/integrations/fastapi/middleware.py @@ -5,23 +5,28 @@ import functools from collections.abc import Callable -from starlette.types import ASGIApp, Message, Receive, Scope, Send +try: + from starlette.types import ASGIApp, Message, Receive, Scope, Send +except ImportError as e: # pragma: no cover + raise ImportError( + "The CodeCarbon FastAPI integration needs FastAPI/Starlette: " + "pip install 'codecarbon[fastapi]'" + ) from e from codecarbon.emissions_tracker import BaseEmissionsTracker from codecarbon.external.logger import logger from codecarbon.integrations.fastapi.attribution import EnergyAttributor, RequestEnergy -from codecarbon.output_methods.emissions_data import EmissionsData def log_request( - energy: RequestEnergy, emissions: EmissionsData | None, status_code: int + energy: RequestEnergy, emissions_kg: float | None, status_code: int ) -> None: """Default ``on_request`` handler; logs via the ``codecarbon`` logger.""" - logger.info( + logger.debug( "CodeCarbon %s: energy=%s kWh emissions=%s kg CO2 status=%s", energy.endpoint, energy.energy_kwh, - getattr(emissions, "emissions", None), + emissions_kg, status_code, ) @@ -29,14 +34,19 @@ def log_request( class CodeCarbonMiddleware: """Attributes a running tracker's energy to each HTTP request. + Add it at module level with ``app.add_middleware(CodeCarbonMiddleware)``. + The tracker is either passed as ``tracker=`` or looked up per request from + ``app.state.codecarbon_tracker``, so it can be created and started in the + app's lifespan. Requests are only recorded while that tracker is running. + A request's number is only known one or more sampling windows *after* its response was sent, so ``on_request`` is called then, from the tracker's scheduler thread. Keep it cheap and non-blocking. Args: app: Inner ASGI application. - tracker: A tracker already started with ``tracker.start()``. - on_request: Callback ``(RequestEnergy, EmissionsData | None, status_code)``. + tracker: Optional tracker; defaults to ``app.state.codecarbon_tracker``. + on_request: Callback ``(RequestEnergy, emissions_kg | None, status_code)``. ``None`` disables reporting. """ @@ -44,29 +54,73 @@ def __init__( self, app: ASGIApp, *, - tracker: BaseEmissionsTracker, + tracker: BaseEmissionsTracker | None = None, on_request: ( - Callable[[RequestEnergy, EmissionsData | None, int], None] | None + Callable[[RequestEnergy, float | None, int], None] | None ) = log_request, ) -> None: self.app = app self.tracker = tracker self.on_request = on_request self.attributor = EnergyAttributor() - self.attributor.reset_window(tracker._total_energy.kWh) - tracker.add_energy_window_observer(self.attributor.on_window) + self._attached: BaseEmissionsTracker | None = None + # kg CO2eq per kWh, refreshed once per sampling window. + self._intensity: float | None = None def close(self) -> None: - """Stop attributing and emit whatever is still in flight.""" - self.tracker.remove_energy_window_observer(self.attributor.on_window) + """Stop attributing and emit whatever is still in flight. + + Called automatically on lifespan shutdown. + """ + if self._attached is not None: + self._attached.remove_energy_window_observer(self._on_window) + self._attached = None self.attributor.close() + def _on_window(self, total_energy_kwh: float) -> None: + # Scheduler thread: one intensity lookup per window, not per request. + try: + self._intensity = self._attached._carbon_intensity_kg_per_kwh() + except Exception: + logger.debug("CodeCarbon: carbon intensity unavailable", exc_info=True) + self.attributor.on_window(total_energy_kwh) + + def _running_tracker(self, scope: Scope) -> BaseEmissionsTracker | None: + tracker = self.tracker + if tracker is None and "app" in scope: + tracker = getattr(scope["app"].state, "codecarbon_tracker", None) + if tracker is None or tracker._start_time is None: + return None + return tracker + async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None: """ASGI entrypoint.""" + if scope["type"] == "lifespan": + + async def lifespan_send(message: Message) -> None: + if message["type"] == "lifespan.shutdown.complete": + self.close() + await send(message) + + await self.app(scope, receive, lifespan_send) + return + if scope["type"] != "http": await self.app(scope, receive, send) return + tracker = self._running_tracker(scope) + if tracker is not self._attached: + # Tracker stopped, replaced or first seen: settle what we hold. + self.close() + if tracker is not None: + self.attributor.reset_window(tracker._total_energy.kWh) + tracker.add_energy_window_observer(self._on_window) + self._attached = tracker + if tracker is None: + await self.app(scope, receive, send) + return + status_code = 500 async def send_wrapper(message: Message) -> None: @@ -88,31 +142,15 @@ async def send_wrapper(message: Message) -> None: def _resolved(self, status_code: int, energy: RequestEnergy) -> None: if self.on_request is None: return - emissions = ( - self.tracker.http_request_emissions(energy.energy_kwh, energy.duration_s) - if energy.energy_kwh is not None + emissions_kg = ( + energy.energy_kwh * self._intensity + if energy.energy_kwh is not None and self._intensity is not None else None ) - self.on_request(energy, emissions, status_code) + self.on_request(energy, emissions_kg, status_code) def _endpoint(scope: Scope) -> str: route = scope.get("route") path = getattr(route, "path", None) or scope.get("path", "") return f"{scope.get('method', '')} {path}".strip() - - -def add_codecarbon_middleware(app, **kwargs) -> None: - """Register :class:`CodeCarbonMiddleware` and expose it on ``app.state``. - - Starlette builds the middleware stack on startup, so the instance that - actually serves requests is only knowable from inside its constructor. - ``app.state.codecarbon_middleware.close()`` on shutdown. - """ - - class _Registered(CodeCarbonMiddleware): - def __init__(self, asgi_app: ASGIApp, **kw) -> None: - super().__init__(asgi_app, **kw) - app.state.codecarbon_middleware = self - - app.add_middleware(_Registered, **kwargs) diff --git a/docs/how-to/examples.md b/docs/how-to/examples.md index 2763cd2c9..a6f607ecf 100644 --- a/docs/how-to/examples.md +++ b/docs/how-to/examples.md @@ -166,30 +166,36 @@ weighted by overlap. Per-request start/stop snapshots cannot be used here: with N requests in flight each one would see the whole machine's delta, so the sum overcounts by roughly N. +Install the extra with `pip install 'codecarbon[fastapi]'`. Add the middleware +at module level (Starlette refuses new middleware once the app has started), +and start and stop the tracker in the lifespan: + ```python from contextlib import asynccontextmanager from fastapi import FastAPI from codecarbon import EmissionsTracker -from codecarbon.integrations.fastapi import add_codecarbon_middleware +from codecarbon.integrations.fastapi import CodeCarbonMiddleware @asynccontextmanager async def lifespan(app: FastAPI): tracker = EmissionsTracker(allow_multiple_runs=True) tracker.start() - add_codecarbon_middleware(app, tracker=tracker) + app.state.codecarbon_tracker = tracker yield tracker.stop() - app.state.codecarbon_middleware.close() app = FastAPI(lifespan=lifespan) +app.add_middleware(CodeCarbonMiddleware) ``` -A request's share is only known one or more sampling windows *after* its -response was sent, so the `on_request(energy, emissions, status_code)` -callback fires then, on the tracker's scheduler thread — keep it cheap. +Requests are only recorded while the tracker is running. Anything still +pending is reported when the app shuts down. A request's share is only known +one or more sampling windows *after* its response was sent, so the +`on_request(energy, emissions_kg, status_code)` callback fires then, on the +tracker's scheduler thread — keep it cheap. The default callback logs at DEBUG. `energy.energy_kwh` is `None` when the request was shorter than the gap between two samples: there is no honest number, and zero would be a lie. diff --git a/pyproject.toml b/pyproject.toml index 7794c5103..8471ec54f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -115,6 +115,9 @@ carbonboard = [ "dash_bootstrap_components > 1.0.0", "fire", ] +fastapi = [ + "fastapi>=0.100", +] # Backwards compatibility alias - will be removed in v4.0.0 viz-legacy = [ "dash", diff --git a/tests/integrations/test_fastapi.py b/tests/integrations/test_fastapi.py index aa7fb4e1a..403949efe 100644 --- a/tests/integrations/test_fastapi.py +++ b/tests/integrations/test_fastapi.py @@ -2,16 +2,18 @@ import threading import time +from contextlib import asynccontextmanager +from types import SimpleNamespace import pytest from fastapi import FastAPI from fastapi.testclient import TestClient -from codecarbon.emissions_tracker import EmissionsTracker +from codecarbon.emissions_tracker import OfflineEmissionsTracker from codecarbon.integrations.fastapi import ( + CodeCarbonMiddleware, EnergyAttributor, RequestEnergy, - add_codecarbon_middleware, ) @@ -110,38 +112,110 @@ def requester(i: int): assert sum(resolved) == pytest.approx(attributor.attributed_kwh) -def test_end_to_end_through_a_real_tracker(): - tracker = EmissionsTracker( - measure_power_secs=0.5, output_methods=[], allow_multiple_runs=True - ) - tracker.start() - seen: list[tuple] = [] - app = FastAPI() - add_codecarbon_middleware( - app, +class _FakeTracker: + """The slice of a tracker the middleware touches; windows closed by hand.""" + + def __init__(self, started: bool = True) -> None: + self._start_time = 0.0 if started else None + self._total_energy = SimpleNamespace(kWh=0.0) + self.observers = [] + + def add_energy_window_observer(self, callback): + self.observers.append(callback) + + def remove_energy_window_observer(self, callback): + self.observers.remove(callback) + + def _carbon_intensity_kg_per_kwh(self): + return 0.5 + + def window(self, kwh: float) -> None: + self._total_energy.kWh += kwh + for callback in tuple(self.observers): + callback(self._total_energy.kWh) + + +def _app(tracker, seen, *, lifespan=None): + app = FastAPI(lifespan=lifespan) + app.add_middleware( + CodeCarbonMiddleware, tracker=tracker, - on_request=lambda energy, emissions, status: seen.append( - (energy, emissions, status) - ), + on_request=lambda energy, kg, status: seen.append((energy, kg, status)), ) @app.get("/work/{n}") def work(n: int): - time.sleep(0.6) # long enough to span a sampling window + time.sleep(0.01) return {"n": n} - try: - with TestClient(app) as client: + @app.get("/boom") + def boom(): + raise RuntimeError("boom") + + return app + + +def test_request_resolves_on_next_window(): + tracker, seen = _FakeTracker(), [] + with TestClient(_app(tracker, seen)) as client: + assert client.get("/work/1").status_code == 200 + assert seen == [] # not known until a window closes + tracker.window(1.0) + energy, kg, status = seen[0] + assert (energy.endpoint, status) == ("GET /work/{n}", 200) + assert energy.energy_kwh == pytest.approx(1.0) + assert kg == pytest.approx(0.5) + + +def test_app_raising_is_reported_as_500(): + tracker, seen = _FakeTracker(), [] + with TestClient(_app(tracker, seen), raise_server_exceptions=False) as client: + assert client.get("/boom").status_code == 500 + tracker.window(1.0) + assert seen[0][2] == 500 + + +def test_close_on_shutdown_emits_pending(): + tracker, seen = _FakeTracker(), [] + with TestClient(_app(tracker, seen)) as client: + client.get("/work/1") + # Lifespan shutdown closed the middleware: no window ever covered it. + assert [(e.energy_kwh, kg) for e, kg, _ in seen] == [(None, None)] + assert tracker.observers == [] + + +def test_tracker_not_started_records_nothing(): + tracker, seen = _FakeTracker(started=False), [] + with TestClient(_app(tracker, seen)) as client: + for _ in range(5): assert client.get("/work/1").status_code == 200 - time.sleep(1.2) # let a window close and resolve the request - finally: + assert seen == [] + assert tracker.observers == [] + + +def test_documented_lifespan_pattern(): + """The pattern from docs/how-to/examples.md, on an offline tracker.""" + seen = [] + + @asynccontextmanager + async def lifespan(app: FastAPI): + tracker = OfflineEmissionsTracker( + country_iso_code="FRA", + measure_power_secs=3600, # windows are closed by hand below + output_methods=[], + allow_multiple_runs=True, + ) + tracker.start() + app.state.codecarbon_tracker = tracker + yield tracker.stop() - app.state.codecarbon_middleware.close() - assert seen, "no request was reported" - energy, emissions, status = seen[0] - assert status == 200 - assert energy.endpoint == "GET /work/{n}" + app = _app(None, seen, lifespan=lifespan) + with TestClient(app) as client: + assert client.get("/work/1").status_code == 200 + app.state.codecarbon_tracker._measure_power_and_energy() + energy, kg, status = seen[0] + + assert (energy.endpoint, status) == ("GET /work/{n}", 200) assert energy.energy_kwh is not None and energy.energy_kwh > 0 - assert emissions.energy_consumed == energy.energy_kwh - assert emissions.duration == pytest.approx(energy.duration_s) + assert kg is not None and kg > 0 diff --git a/uv.lock b/uv.lock index a787e5864..84f52c817 100644 --- a/uv.lock +++ b/uv.lock @@ -568,6 +568,9 @@ carbonboard = [ { name = "dash-bootstrap-components" }, { name = "fire" }, ] +fastapi = [ + { name = "fastapi" }, +] viz-legacy = [ { name = "dash" }, { name = "dash-bootstrap-components" }, @@ -614,6 +617,7 @@ requires-dist = [ { name = "dash", marker = "extra == 'viz-legacy'" }, { name = "dash-bootstrap-components", marker = "extra == 'carbonboard'", specifier = ">1.0.0" }, { name = "dash-bootstrap-components", marker = "extra == 'viz-legacy'", specifier = ">1.0.0" }, + { name = "fastapi", marker = "extra == 'fastapi'", specifier = ">=0.100" }, { name = "fire", marker = "extra == 'carbonboard'" }, { name = "fire", marker = "extra == 'viz-legacy'" }, { name = "joserfc", specifier = ">=1.0.0" }, @@ -631,7 +635,7 @@ requires-dist = [ { name = "rich" }, { name = "typer" }, ] -provides-extras = ["carbonboard", "viz-legacy"] +provides-extras = ["carbonboard", "fastapi", "viz-legacy"] [package.metadata.requires-dev] dev = [ @@ -929,7 +933,7 @@ name = "exceptiongroup" version = "1.3.1" source = { registry = "https://pypi.org/simple" } dependencies = [ - { name = "typing-extensions" }, + { name = "typing-extensions", marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/50/79/66800aadf48771f6b62f7eb014e352e5d06856655206165d775e675a02c9/exceptiongroup-1.3.1.tar.gz", hash = "sha256:8b412432c6055b0b7d14c310000ae93352ed6754f70fa8f7c34141f91c4e3219", size = 30371, upload-time = "2025-11-21T23:01:54.787Z" } wheels = [ @@ -2067,10 +2071,10 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } }, - { name = "python-dateutil" }, - { name = "pytz" }, - { name = "tzdata" }, + { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, + { name = "python-dateutil", marker = "python_full_version < '3.11'" }, + { name = "pytz", marker = "python_full_version < '3.11'" }, + { name = "tzdata", marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/33/01/d40b85317f86cf08d853a4f495195c73815fdf205eef3993821720274518/pandas-2.3.3.tar.gz", hash = "sha256:e05e1af93b977f7eafa636d043f9f94c7ee3ac81af99c13508215942e64c993b", size = 4495223, upload-time = "2025-09-29T23:34:51.853Z" } wheels = [ @@ -2142,10 +2146,10 @@ resolution-markers = [ "python_full_version == '3.11.*' and sys_platform != 'emscripten' and sys_platform != 'win32'", ] dependencies = [ - { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" }, + { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" }, { name = "numpy", version = "2.5.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" }, - { name = "python-dateutil" }, - { name = "tzdata", marker = "sys_platform == 'emscripten' or sys_platform == 'win32'" }, + { name = "python-dateutil", marker = "python_full_version >= '3.11'" }, + { name = "tzdata", marker = "(python_full_version >= '3.11' and sys_platform == 'emscripten') or (python_full_version >= '3.11' and sys_platform == 'win32')" }, ] sdist = { url = "https://files.pythonhosted.org/packages/be/4f/5f3422a2afec5ffc46308b79e53291365a93748b498ac2e58bead0197916/pandas-3.0.5.tar.gz", hash = "sha256:dca3734d6ab7c906e6730f0788b0a1dbb9f2467731f9711f77995c8e9d62d712", size = 4658219, upload-time = "2026-07-22T22:19:28.819Z" } wheels = [ @@ -3268,10 +3272,10 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "joblib" }, - { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } }, - { name = "scipy", version = "1.15.3", source = { registry = "https://pypi.org/simple" } }, - { name = "threadpoolctl" }, + { name = "joblib", marker = "python_full_version < '3.11'" }, + { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, + { name = "scipy", version = "1.15.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, + { name = "threadpoolctl", marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/98/c2/a7855e41c9d285dfe86dc50b250978105dce513d6e459ea66a6aeb0e1e0c/scikit_learn-1.7.2.tar.gz", hash = "sha256:20e9e49ecd130598f1ca38a1d85090e1a600147b9c02fa6f15d69cb53d968fda", size = 7193136, upload-time = "2025-09-09T08:21:29.075Z" } wheels = [ @@ -3326,13 +3330,13 @@ resolution-markers = [ "python_full_version == '3.11.*' and sys_platform != 'emscripten' and sys_platform != 'win32'", ] dependencies = [ - { name = "joblib" }, - { name = "narwhals" }, - { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" }, + { name = "joblib", marker = "python_full_version >= '3.11'" }, + { name = "narwhals", marker = "python_full_version >= '3.11'" }, + { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" }, { name = "numpy", version = "2.5.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" }, - { name = "scipy", version = "1.17.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" }, + { name = "scipy", version = "1.17.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" }, { name = "scipy", version = "1.18.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" }, - { name = "threadpoolctl" }, + { name = "threadpoolctl", marker = "python_full_version >= '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/fa/6f/37092bdb25f712817231799fc5674d8e704066a8a70c1d2d40517e18b4ab/scikit_learn-1.9.0.tar.gz", hash = "sha256:8833266989d3a5110178a9fae30783675460724d0e1efb13b14901d2c660c557", size = 7750767, upload-time = "2026-06-02T11:54:32.706Z" } wheels = [ @@ -3376,7 +3380,7 @@ resolution-markers = [ "python_full_version < '3.11'", ] dependencies = [ - { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } }, + { name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/0f/37/6964b830433e654ec7485e45a00fc9a27cf868d622838f6b6d9c5ec0d532/scipy-1.15.3.tar.gz", hash = "sha256:eae3cf522bc7df64b42cad3925c876e1b0b6c35c1337c93e12c0f366f55b0eaf", size = 59419214, upload-time = "2025-05-08T16:13:05.955Z" } wheels = [ @@ -3437,7 +3441,7 @@ resolution-markers = [ "python_full_version == '3.11.*' and sys_platform != 'emscripten' and sys_platform != 'win32'", ] dependencies = [ - { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" } }, + { name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/7a/97/5a3609c4f8d58b039179648e62dd220f89864f56f7357f5d4f45c29eb2cc/scipy-1.17.1.tar.gz", hash = "sha256:95d8e012d8cb8816c226aef832200b1d45109ed4464303e997c5b13122b297c0", size = 30573822, upload-time = "2026-02-23T00:26:24.851Z" } wheels = [ @@ -3519,7 +3523,7 @@ resolution-markers = [ "python_full_version >= '3.12' and python_full_version < '3.14' and sys_platform != 'emscripten' and sys_platform != 'win32'", ] dependencies = [ - { name = "numpy", version = "2.5.3", source = { registry = "https://pypi.org/simple" } }, + { name = "numpy", version = "2.5.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" }, ] sdist = { url = "https://files.pythonhosted.org/packages/7e/74/66de6258867beb2ef08f35f9f2ac017a52cacd5081714d239ff1a442d458/scipy-1.18.1.tar.gz", hash = "sha256:52c4b7422442aba924d03ad4019852b08a92e64ea187b933135687bfe2747307", size = 30781235, upload-time = "2026-08-21T23:28:50.599Z" } wheels = [ From 4732135f8544cd76944ae99dd0d4a8d464eea870 Mon Sep 17 00:00:00 2001 From: David Berenstein Date: Sun, 27 Sep 2026 16:12:54 +0900 Subject: [PATCH 3/4] docs: explain FastAPI per-request energy is an estimated share --- docs/how-to/examples.md | 13 +++++++++++-- 1 file changed, 11 insertions(+), 2 deletions(-) diff --git a/docs/how-to/examples.md b/docs/how-to/examples.md index a6f607ecf..60b2a2122 100644 --- a/docs/how-to/examples.md +++ b/docs/how-to/examples.md @@ -197,5 +197,14 @@ pending is reported when the app shuts down. A request's share is only known one or more sampling windows *after* its response was sent, so the `on_request(energy, emissions_kg, status_code)` callback fires then, on the tracker's scheduler thread — keep it cheap. The default callback logs at DEBUG. -`energy.energy_kwh` is `None` when the request was shorter than the gap -between two samples: there is no honest number, and zero would be a lie. +`energy.energy_kwh` is `None` only when the request never overlapped a +completed sampling window, which in practice means it was still pending when +the tracker stopped. + +`energy_kwh` is an estimated share, not a measurement of the request. The +whole machine's energy for a window, idle power included, is split across the +requests in flight by how long each overlapped the window. Time spent waiting +on I/O counts the same as time spent computing, and a lone short request in an +otherwise idle window receives that window's full energy. Sum the values per +route over many requests rather than reading a single one, and lower +`measure_power_secs` for finer-grained windows. From edd19133baa4312e02c222b26d79a04b0f27ffee Mon Sep 17 00:00:00 2001 From: David Berenstein Date: Mon, 28 Sep 2026 12:26:40 +0900 Subject: [PATCH 4/4] fix(fastapi): refresh carbon intensity at the tracker's API cadence, document multi-worker overcount --- codecarbon/integrations/fastapi/middleware.py | 22 ++++++++++++++----- docs/how-to/examples.md | 9 ++++++++ tests/integrations/test_fastapi.py | 17 +++++++++++++- 3 files changed, 41 insertions(+), 7 deletions(-) diff --git a/codecarbon/integrations/fastapi/middleware.py b/codecarbon/integrations/fastapi/middleware.py index e44ccc527..bfbd79d1e 100644 --- a/codecarbon/integrations/fastapi/middleware.py +++ b/codecarbon/integrations/fastapi/middleware.py @@ -64,8 +64,10 @@ def __init__( self.on_request = on_request self.attributor = EnergyAttributor() self._attached: BaseEmissionsTracker | None = None - # kg CO2eq per kWh, refreshed once per sampling window. + # kg CO2eq per kWh, refreshed at most once per `api_call_interval` + # windows -- the same cadence the tracker itself uses to refresh. self._intensity: float | None = None + self._windows_since_refresh = 0 def close(self) -> None: """Stop attributing and emit whatever is still in flight. @@ -75,14 +77,22 @@ def close(self) -> None: if self._attached is not None: self._attached.remove_energy_window_observer(self._on_window) self._attached = None + self._windows_since_refresh = 0 self.attributor.close() def _on_window(self, total_energy_kwh: float) -> None: - # Scheduler thread: one intensity lookup per window, not per request. - try: - self._intensity = self._attached._carbon_intensity_kg_per_kwh() - except Exception: - logger.debug("CodeCarbon: carbon intensity unavailable", exc_info=True) + # Scheduler thread: refresh intensity no more often than the tracker + # itself does (its `api_call_interval`), since an Electricity Maps + # token turns this into an HTTP call. + interval = getattr(self._attached, "_api_call_interval", 1) or 1 + if self._intensity is None or self._windows_since_refresh >= interval: + try: + self._intensity = self._attached._carbon_intensity_kg_per_kwh() + except Exception: + logger.debug("CodeCarbon: carbon intensity unavailable", exc_info=True) + self._windows_since_refresh = 0 + else: + self._windows_since_refresh += 1 self.attributor.on_window(total_energy_kwh) def _running_tracker(self, scope: Scope) -> BaseEmissionsTracker | None: diff --git a/docs/how-to/examples.md b/docs/how-to/examples.md index 60b2a2122..0fd93bb89 100644 --- a/docs/how-to/examples.md +++ b/docs/how-to/examples.md @@ -208,3 +208,12 @@ on I/O counts the same as time spent computing, and a lone short request in an otherwise idle window receives that window's full energy. Sum the values per route over many requests rather than reading a single one, and lower `measure_power_secs` for finer-grained windows. + +With `uvicorn --workers N` (or any multi-process server), each worker gets +its own tracker in its own process, and by default a tracker measures the +*whole machine*. Summing per-request energy across all workers then +overcounts by roughly N×, the same overcounting the middleware avoids within +a single process. Run a single worker per machine when you need per-request +numbers. `tracking_mode="process"` only helps when CPU power is estimated +from CPU load: RAPL, powermetrics and NVML counters are machine-wide, so each +worker still sees the whole machine's energy. diff --git a/tests/integrations/test_fastapi.py b/tests/integrations/test_fastapi.py index 403949efe..758a2f213 100644 --- a/tests/integrations/test_fastapi.py +++ b/tests/integrations/test_fastapi.py @@ -115,10 +115,12 @@ def requester(i: int): class _FakeTracker: """The slice of a tracker the middleware touches; windows closed by hand.""" - def __init__(self, started: bool = True) -> None: + def __init__(self, started: bool = True, api_call_interval: int = 8) -> None: self._start_time = 0.0 if started else None self._total_energy = SimpleNamespace(kWh=0.0) + self._api_call_interval = api_call_interval self.observers = [] + self.intensity_calls = 0 def add_energy_window_observer(self, callback): self.observers.append(callback) @@ -127,6 +129,7 @@ def remove_energy_window_observer(self, callback): self.observers.remove(callback) def _carbon_intensity_kg_per_kwh(self): + self.intensity_calls += 1 return 0.5 def window(self, kwh: float) -> None: @@ -167,6 +170,18 @@ def test_request_resolves_on_next_window(): assert kg == pytest.approx(0.5) +def test_intensity_refreshed_once_per_api_call_interval(): + # An Electricity Maps token turns _carbon_intensity_kg_per_kwh into an + # HTTP call; it must not fire on every sampling window, only as often as + # the tracker itself refreshes (its api_call_interval). + tracker, seen = _FakeTracker(api_call_interval=4), [] + with TestClient(_app(tracker, seen)) as client: + client.get("/work/1") # attaches the middleware to the tracker + for _ in range(12): + tracker.window(0.1) + assert tracker.intensity_calls == pytest.approx(3, abs=1) + + def test_app_raising_is_reported_as_500(): tracker, seen = _FakeTracker(), [] with TestClient(_app(tracker, seen), raise_server_exceptions=False) as client: