diff --git a/codecarbon/emissions_tracker.py b/codecarbon/emissions_tracker.py index 4bd0ea99d..f14a5ffdd 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._scheduler_paused_by_task = False self._hardware = [] self._hardware_initialized = False @@ -1004,6 +1005,46 @@ 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: + # 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 _carbon_intensity_kg_per_kwh(self) -> float: + """Current carbon intensity, kg CO2eq per kWh, without touching run totals. + + Used by the FastAPI integration to convert each request's energy share, + once per sampling window, off the tracker's accumulated state. + """ + 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: """ Prepare the emissions data to be sent to the API or written to a file. @@ -1304,6 +1345,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..dd1969761 --- /dev/null +++ b/codecarbon/integrations/fastapi/__init__.py @@ -0,0 +1,10 @@ +"""FastAPI integration: per-request energy attribution middleware.""" + +from codecarbon.integrations.fastapi.attribution import EnergyAttributor, RequestEnergy +from codecarbon.integrations.fastapi.middleware import CodeCarbonMiddleware + +__all__ = [ + "CodeCarbonMiddleware", + "EnergyAttributor", + "RequestEnergy", +] 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..bfbd79d1e --- /dev/null +++ b/codecarbon/integrations/fastapi/middleware.py @@ -0,0 +1,166 @@ +"""FastAPI/Starlette middleware for per-request emissions attribution.""" + +from __future__ import annotations + +import functools +from collections.abc import Callable + +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 + + +def log_request( + energy: RequestEnergy, emissions_kg: float | None, status_code: int +) -> None: + """Default ``on_request`` handler; logs via the ``codecarbon`` logger.""" + logger.debug( + "CodeCarbon %s: energy=%s kWh emissions=%s kg CO2 status=%s", + energy.endpoint, + energy.energy_kwh, + emissions_kg, + status_code, + ) + + +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: Optional tracker; defaults to ``app.state.codecarbon_tracker``. + on_request: Callback ``(RequestEnergy, emissions_kg | None, status_code)``. + ``None`` disables reporting. + """ + + def __init__( + self, + app: ASGIApp, + *, + tracker: BaseEmissionsTracker | None = None, + on_request: ( + 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._attached: BaseEmissionsTracker | None = None + # 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. + + Called automatically on lifespan shutdown. + """ + 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: 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: + 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: + 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_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_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() diff --git a/docs/how-to/examples.md b/docs/how-to/examples.md index 06f084aca..0fd93bb89 100644 --- a/docs/how-to/examples.md +++ b/docs/how-to/examples.md @@ -157,3 +157,63 @@ 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. + +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 CodeCarbonMiddleware + + +@asynccontextmanager +async def lifespan(app: FastAPI): + tracker = EmissionsTracker(allow_multiple_runs=True) + tracker.start() + app.state.codecarbon_tracker = tracker + yield + tracker.stop() + + +app = FastAPI(lifespan=lifespan) +app.add_middleware(CodeCarbonMiddleware) +``` + +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` 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. + +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/pyproject.toml b/pyproject.toml index 9f5d2fef6..8471ec54f 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", @@ -113,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 new file mode 100644 index 000000000..758a2f213 --- /dev/null +++ b/tests/integrations/test_fastapi.py @@ -0,0 +1,236 @@ +"""Tests for the FastAPI per-request energy attribution.""" + +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 OfflineEmissionsTracker +from codecarbon.integrations.fastapi import ( + CodeCarbonMiddleware, + EnergyAttributor, + RequestEnergy, +) + + +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) + + +class _FakeTracker: + """The slice of a tracker the middleware touches; windows closed by hand.""" + + 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) + + 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: + 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, kg, status: seen.append((energy, kg, status)), + ) + + @app.get("/work/{n}") + def work(n: int): + time.sleep(0.01) + return {"n": n} + + @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_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: + 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 + 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 = _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 kg is not None and kg > 0 diff --git a/uv.lock b/uv.lock index e9cec04bf..84f52c817 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" @@ -554,6 +568,9 @@ carbonboard = [ { name = "dash-bootstrap-components" }, { name = "fire" }, ] +fastapi = [ + { name = "fastapi" }, +] viz-legacy = [ { name = "dash" }, { name = "dash-bootstrap-components" }, @@ -564,6 +581,8 @@ viz-legacy = [ dev = [ { name = "black" }, { name = "bumpver" }, + { name = "fastapi" }, + { name = "httpx" }, { name = "jsonschema" }, { name = "logfire" }, { name = "mktestdocs" }, @@ -598,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" }, @@ -615,12 +635,14 @@ requires-dist = [ { name = "rich" }, { name = "typer" }, ] -provides-extras = ["carbonboard", "viz-legacy"] +provides-extras = ["carbonboard", "fastapi", "viz-legacy"] [package.metadata.requires-dev] dev = [ { name = "black" }, { name = "bumpver" }, + { name = "fastapi", specifier = ">=0.100" }, + { name = "httpx" }, { name = "jsonschema" }, { name = "logfire", specifier = ">=1.0.1" }, { name = "mktestdocs" }, @@ -911,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 = [ @@ -927,6 +949,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.6" @@ -998,6 +1036,43 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/41/63/e876e789525063c840ccfa8857febdabd6523bcef9ce7eb979b9305ea895/griffelib-2.3.0-py3-none-any.whl", hash = "sha256:1b8f9cd525681c26b1d6d574faa1371651e8459ca51d209684f50b8096ae06e0", size = 169423, upload-time = "2026-09-04T15:08:12.956Z" }, ] +[[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" @@ -1996,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 = [ @@ -2071,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 = [ @@ -3197,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 = [ @@ -3255,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 = [ @@ -3305,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 = [ @@ -3366,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 = [ @@ -3448,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 = [ @@ -3550,6 +3625,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/eb/dc/ad025c1ee131eba60c69f4dd5779b18fcf1e6b21a343e2162a84d5d133c7/soupsieve-2.9.2-py3-none-any.whl", hash = "sha256:8089a26fd974ca7a1f30276d3d8492ab266ab15af581642dfe8aa162e0c1c823", size = 37370, upload-time = "2026-08-07T00:57:23.524Z" }, ] +[[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"