diff --git a/codecarbon/emissions_tracker.py b/codecarbon/emissions_tracker.py index f14a5ffdd..a09ca55c7 100644 --- a/codecarbon/emissions_tracker.py +++ b/codecarbon/emissions_tracker.py @@ -23,7 +23,13 @@ from codecarbon.core.config import get_hierarchical_config, normalize_gpu_ids from codecarbon.core.units import Energy, Power, Time, Water from codecarbon.core.util import count_cpus, count_physical_cpus, suppress -from codecarbon.external.hardware import CPU, GPU, AppleSiliconChip +from codecarbon.external.hardware import ( + CONSUMPTION_PERCENTAGE_CONSTANT, + CPU, + GPU, + MODE_CPU_LOAD, + AppleSiliconChip, +) from codecarbon.external.logger import logger, set_logger_format, set_logger_level from codecarbon.external.ram import RAM from codecarbon.external.scheduler import PeriodicScheduler @@ -55,6 +61,39 @@ _sentinel = object() +#: CPU modes that read an energy or power counter from the hardware. +_MEASURED_CPU_MODES = ("intel_rapl", "windows_emi", "intel_power_gadget") + + +@dataclasses.dataclass(frozen=True) +class WindowSample: + """The tracker's state at the end of one sampling window. + + Energies are cumulative since the tracker started, in kWh, PUE included. + Quality is ``"measured"`` for a hardware counter (RAPL, EMI, Power Gadget, + powermetrics, NVML), ``"modeled"`` for CPU load mode and ``"none"`` for a + constant TDP guess. + """ + + #: ``time.perf_counter()`` when the sample was taken. + timestamp: float + total_kwh: float + cpu_kwh: float + gpu_kwh: float + ram_kwh: float + cpu_quality: str + #: ``None`` when no GPU is tracked. + gpu_quality: Optional[str] + #: CPU idle power in W when the power model fixes it (load and constant + #: modes), ``None`` when it has to be estimated from the measurements. + cpu_idle_w: Optional[float] + #: Extra CPU power per busy logical CPU in W (J per CPU-second) when the + #: power model fixes it, ``None`` when it has to be fitted. + cpu_w_per_busy_cpu: Optional[float] + #: The CPU energy already covers only this process (load mode with + #: ``tracking_mode="process"``), not the whole machine. + cpu_per_process: bool + class BaseEmissionsTracker(ABC): """ @@ -296,7 +335,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._window_observers: List[Callable[[WindowSample], None]] = [] self._scheduler_paused_by_task = False self._hardware = [] self._hardware_initialized = False @@ -1005,8 +1044,10 @@ 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. + def add_energy_window_observer( + self, callback: Callable[[WindowSample], None] + ) -> None: + """Call ``callback(sample)`` 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 @@ -1014,22 +1055,83 @@ def add_energy_window_observer(self, callback: Callable[[float], None]) -> None: that were in flight during it. Args: - callback: Receives the tracker's cumulative energy in kWh. + callback: Receives a :class:`WindowSample` with the tracker's + cumulative energy per component. """ self._window_observers.append(callback) - def remove_energy_window_observer(self, callback: Callable[[float], None]) -> None: + def remove_energy_window_observer( + self, callback: Callable[[WindowSample], None] + ) -> None: """Remove a callback registered with :meth:`add_energy_window_observer`.""" if callback in self._window_observers: self._window_observers.remove(callback) + def _window_sample(self) -> WindowSample: + """Snapshot of the cumulative energies and how they were obtained.""" + cpu_quality, gpu_quality = "none", None + cpu_idle_w: Optional[float] = None + cpu_w_per_busy_cpu: Optional[float] = None + cpu_per_process = False + for hardware in self._hardware: + if isinstance(hardware, CPU): + if hardware._mode in _MEASURED_CPU_MODES: + cpu_quality = "measured" + elif hardware._mode == MODE_CPU_LOAD: + cpu_quality = "modeled" + cpu_per_process = hardware._tracking_mode == "process" + # Machine load mode draws tdp * (0.1 + 0.9 * load^3), so + # idle is exactly 0.1 * TDP. Process mode has no floor. + cpu_idle_w = 0.0 if cpu_per_process else 0.1 * hardware._tdp + cpu_idle_w *= self._pue + # Process mode is linear in CPU time: TDP / CPUs per busy + # CPU. Machine mode is cubic in load; its chord from idle + # to full load, 0.9 * TDP / CPUs, is the linear stand-in. + if cpu_per_process: + slope, cpus = hardware._tdp, hardware._cpu_count + else: + slope, cpus = 0.9 * hardware._tdp, psutil.cpu_count() + cpu_w_per_busy_cpu = slope * self._pue / max(cpus or 1, 1) + elif hardware._mode == "constant": + # Constant mode never moves: all of it is idle power. + cpu_idle_w = ( + hardware._tdp * CONSUMPTION_PERCENTAGE_CONSTANT * self._pue + ) + cpu_w_per_busy_cpu = 0.0 + elif isinstance(hardware, AppleSiliconChip): + if hardware.chip_part == "CPU": + cpu_quality = "measured" + elif hardware.chip_part == "GPU": + gpu_quality = "measured" + elif isinstance(hardware, GPU): + gpu_quality = "measured" + return WindowSample( + timestamp=self._last_measured_time, + total_kwh=self._total_energy.kWh, + cpu_kwh=self._total_cpu_energy.kWh, + gpu_kwh=self._total_gpu_energy.kWh, + ram_kwh=self._total_ram_energy.kWh, + cpu_quality=cpu_quality, + gpu_quality=gpu_quality, + cpu_idle_w=cpu_idle_w, + cpu_w_per_busy_cpu=cpu_w_per_busy_cpu, + cpu_per_process=cpu_per_process, + ) + 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") + if not self._window_observers: + return + # Observers are integrations: nothing they touch may break measurement. + try: + sample = self._window_sample() + # Copy: an observer may be removed from another thread mid-iteration. + for callback in tuple(self._window_observers): + try: + callback(sample) + except Exception: + logger.exception("CodeCarbon energy window observer failed") + except Exception: + logger.exception("CodeCarbon energy window sample failed") def _carbon_intensity_kg_per_kwh(self) -> float: """Current carbon intensity, kg CO2eq per kWh, without touching run totals. diff --git a/codecarbon/integrations/fastapi/attribution.py b/codecarbon/integrations/fastapi/attribution.py index eb17d30bc..b862c17b1 100644 --- a/codecarbon/integrations/fastapi/attribution.py +++ b/codecarbon/integrations/fastapi/attribution.py @@ -1,13 +1,27 @@ """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:: +Each completed sampling window ``(t_prev, t_now, dE)`` is first split per +component into idle and dynamic energy:: - attributed_kwh + unattributed_kwh == settled_kwh + dynamic_c = max(dE_c - P_idle_c * width, 0) -exactly, after every window. That is the property the tests pin down. +Idle energy (and all RAM energy) goes to ``idle_kwh``: it would have been +drawn with no request at all, so charging it to whichever request happened to +be in flight is wrong. Of the dynamic CPU energy, only this process's share +of the machine's busy CPU time is kept; the rest goes to +``other_processes_kwh``. Each request in flight is charged the CPU time its +meter saw in the window times the cost of a CPU-second; our CPU time no meter +claimed goes to ``process_unattributed_kwh`` at the same cost. One cap scales +both down when all of our CPU time at that cost exceeds what was kept, and +whatever the cost leaves over joins ``other_processes_kwh``. Dynamic GPU energy has no per-request signal +and is split by overlap with the window. Dynamic energy of windows with +nothing in flight goes to ``unattributed_kwh``. The invariant is:: + + attributed_kwh + idle_kwh + other_processes_kwh + + process_unattributed_kwh + unattributed_kwh == settled_kwh + +up to float rounding, 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 @@ -18,28 +32,157 @@ import threading import time +from collections import deque from collections.abc import Callable from dataclasses import dataclass from typing import Any +import psutil + +from codecarbon.emissions_tracker import WindowSample from codecarbon.external.logger import logger +#: Watt-seconds per kWh. +_WS_PER_KWH = 3.6e6 +#: Weakest first: a request's quality is that of its weakest input. +_QUALITIES = ("none", "modeled", "measured") + @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. + ``energy_kwh`` and ``gpu_kwh`` are ``None`` when the request never covered + a completed sampling window: there is no honest number, and zero would be + a lie. """ endpoint: str + #: CPU energy above idle, charged by the CPU time the request used. energy_kwh: float | None + #: GPU energy above idle, charged by wall-clock overlap with the window. + gpu_kwh: float | None duration_s: float - #: Completed sampling windows this request overlapped. + #: Completed sampling windows this request was in flight for. windows: int - #: Mean number of requests it competed against, window-weighted. + #: Mean number of requests in flight with it, per window. mean_concurrency: float | None + #: CPU time the meter saw this request use, in seconds. + cpu_seconds: float + #: ``"cpu_time"``, ``"wall"`` (GPU energy only) or ``"mixed"`` (both). + attribution_method: str + #: ``"measured"`` (RAPL, powermetrics, NVML...), ``"modeled"`` (CPU load + #: mode) or ``"none"`` (constant TDP, or no window), for the weakest + #: component the request was charged from. + quality: str + + +def _thread_clock() -> int | None: + """The calling thread's CPU clock, readable from other threads (Linux).""" + if not hasattr(time, "pthread_getcpuclockid"): + return None + try: + return time.pthread_getcpuclockid(threading.get_ident()) + except OSError: + return None + + +class _Meter: + """CPU time used by one request, in ns. + + ``ns`` is the event-loop time. It is written only by the loop thread + driving the request, so it needs no lock, with or without the GIL. + Worker-thread calls (:meth:`run_in_thread`) go through ``_lock``. While + one runs, :meth:`total_ns` reads its CPU clock from the outside where the + platform allows it (Linux), so a long sync call is charged window by + window. Elsewhere its CPU time only shows up once it returns. + """ + + __slots__ = ("ns", "_lock", "_thread_ns", "_running") + + def __init__(self) -> None: + self.ns = 0 + self._lock = threading.Lock() + self._thread_ns = 0 + #: Worker thread id -> (its CPU clock or None, thread time at start). + self._running: dict[int, tuple[int | None, int]] = {} + + def run_in_thread(self, func: Callable[..., Any], *args: Any) -> Any: + """Call ``func(*args)`` on this worker thread, adding its CPU time.""" + ident = threading.get_ident() + clock = _thread_clock() + start = time.thread_time_ns() + with self._lock: + self._running[ident] = (clock, start) + try: + return func(*args) + finally: + # Read under the lock, or a total_ns() between the read and the + # lock could count more of this call than we add, and go backwards. + with self._lock: + del self._running[ident] + self._thread_ns += time.thread_time_ns() - start + + def total_ns(self) -> int: + """CPU time so far, including worker calls still running.""" + with self._lock: + total = self.ns + self._thread_ns + for clock, start in self._running.values(): + if clock is not None: + try: + total += max(time.clock_gettime_ns(clock) - start, 0) + except OSError: + pass + return total + + +class _Metered: + """Awaitable that drives the awaitable ``coro`` and meters it. + + Adds ``time.thread_time_ns()`` spent inside each ``send``/``throw`` of its + ``__await__`` iterator to ``meter``. ASGI only promises an awaitable, so + that iterator, not ``coro`` itself, is what gets driven. Work the coroutine hands to other threads or to + child tasks is not seen. + """ + + __slots__ = ("_coro", "_meter") + + def __init__(self, coro: Any, meter: _Meter) -> None: + self._coro = coro + self._meter = meter + + def __await__(self): + coro, meter, clock = self._coro.__await__(), self._meter, time.thread_time_ns + value: Any = None + error: BaseException | None = None + while True: + start = clock() + try: + if error is None: + yielded = coro.send(value) + else: + yielded = coro.throw(error) + except StopIteration as stop: + meter.ns += clock() - start + return stop.value + except BaseException: + meter.ns += clock() - start + raise + meter.ns += clock() - start + value = error = None + try: + value = yield yielded + except GeneratorExit: + start = clock() + try: + close = getattr(coro, "close", None) + if close is not None: + close() + finally: + meter.ns += clock() - start + raise + except BaseException as exc: + error = exc @dataclass @@ -48,45 +191,169 @@ class _InFlight: endpoint: str start: float + meter: _Meter | None = None + #: ``meter.total_ns()`` already charged in earlier windows. + cpu_seen_ns: int = 0 end: float | None = None energy: float = 0.0 + gpu_energy: float = 0.0 + quality: str | None = None windows: int = 0 concurrency_sum: float = 0.0 #: Called with the :class:`RequestEnergy` when this request resolves. on_resolved: Callable[[RequestEnergy], None] | None = None +def _cpu_times() -> tuple[float, float, float]: + """``(process CPU s, machine busy CPU s, machine total CPU s)``, cumulative. + + Busy excludes idle, iowait and steal (time the hypervisor gave to another + guest). ``guest`` and ``guest_nice`` are already counted in ``user`` and + ``nice`` on Linux, so they are taken out of the total, as psutil does. + """ + t = psutil.cpu_times() + total = sum(t) - getattr(t, "guest", 0.0) - getattr(t, "guest_nice", 0.0) + busy = total - t.idle - getattr(t, "iowait", 0.0) - getattr(t, "steal", 0.0) + return time.process_time(), busy, total + + +class _IdleEstimator: + """Idle power of one component in W, and its cost per busy CPU. + + Idle power is the lower of the rolling minimum of window power over + ``horizon_s`` and, once 20 windows with a CPU load are known and the fit + is good (R² > 0.8), the intercept of ``power = a + b * busy_cpus``. The + minimum alone overestimates idle on a server that is never idle; the + intercept alone is noise when load barely varies. The fit's slope ``b`` + (W per busy CPU, i.e. J per CPU-second) is kept in :attr:`slope`. + """ + + def __init__(self, horizon_s: float) -> None: + self._horizon_s = horizon_s + # (t, power) with increasing power: the front is the window minimum. + self._mins: deque[tuple[float, float]] = deque() + self._points: deque[tuple[float, float, float]] = deque() + #: Latest estimate, ``None`` before the first window. + self.watts: float | None = None + #: Slope of a good fit, ``None`` without one. + self.slope: float | None = None + + def update(self, t: float, power: float, busy_cpus: float | None = None) -> float: + cutoff = t - self._horizon_s + while self._mins and self._mins[0][0] < cutoff: + self._mins.popleft() + while self._mins and self._mins[-1][1] >= power: + self._mins.pop() + self._mins.append((t, power)) + estimate = self._mins[0][1] + self.slope = None # no load reading, no current fit to charge by + if busy_cpus is not None: + self._points.append((t, busy_cpus, power)) + while self._points[0][0] < cutoff: + self._points.popleft() + fit = _fit(self._points) + self.slope = fit[1] if fit is not None and fit[1] > 0 else None + if fit is not None: + estimate = min(estimate, max(fit[0], 0.0)) + self.watts = estimate + return estimate + + +def _fit(points: deque[tuple[float, float, float]]) -> tuple[float, float] | None: + """Least-squares ``(intercept, slope)`` of power on load, ``None`` if unreliable.""" + # ponytail: full O(n) refit per window, n = horizon / measure_power_secs + # (3600 at 1 s, about 1 ms). Keep running sums if that ever shows up. + n = len(points) + if n < 20: + return None + mean_x = sum(x for _, x, _ in points) / n + mean_p = sum(p for _, _, p in points) / n + sxx = sum((x - mean_x) ** 2 for _, x, _ in points) + syy = sum((p - mean_p) ** 2 for _, _, p in points) + if sxx <= 0 or syy <= 0: + return None + sxy = sum((x - mean_x) * (p - mean_p) for _, x, p in points) + if sxy * sxy / (sxx * syy) <= 0.8: + return None + slope = sxy / sxx + return mean_p - slope * mean_x, slope + + +def _dynamic(delta_kwh: float, idle_w: float, width: float) -> float: + """Energy above idle in one window, within ``[0, delta_kwh]``.""" + return min(max(delta_kwh - idle_w * width / _WS_PER_KWH, 0.0), delta_kwh) + + class EnergyAttributor: - """Splits each sampling window's energy across the requests in flight.""" + """Splits each sampling window's dynamic energy across the requests in flight. - def __init__(self) -> None: + Args: + clock: Must be the clock of :attr:`WindowSample.timestamp`. + cpu_times: Returns cumulative ``(process, busy, total)`` CPU seconds; + replaceable for tests. + idle_horizon_s: How far back the idle-power estimate looks. + """ + + def __init__( + self, + *, + clock: Callable[[], float] = time.perf_counter, + cpu_times: Callable[[], tuple[float, float, float]] = _cpu_times, + idle_horizon_s: float = 3600.0, + ) -> None: + self._clock = clock + self._cpu_times = cpu_times 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. + #: Energy the machine would have drawn with no load (idle power, and + #: all of RAM), kWh. + self.idle_kwh = 0.0 + #: Dynamic CPU energy of the machine's other processes, kWh. + self.other_processes_kwh = 0.0 + #: Our CPU dynamic energy that no request's meter claimed: work in + #: threads, child tasks, the server itself, kWh. + self.process_unattributed_kwh = 0.0 + #: Dynamic 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. + #: Energy taken in from closed windows. The buckets add up to it, up + #: to float rounding, after every window; it is below the tracker's + #: run total by whatever skipped windows dropped (``windows_skipped``) + #: plus the final unsampled partial window. The anchor sample is read + #: on the request path while the scheduler thread may be measuring, + #: so it can be taken mid-update: the first window is then split + #: slightly off, but everything later is unaffected. self.settled_kwh = 0.0 self.windows_settled = 0 - #: Windows where the energy counter went backwards (RAPL wrap/reset). + #: Windows dropped unsplit: an energy counter went backwards (RAPL + #: wrap/reset), or energy arrived in a window with no width. self.windows_skipped = 0 - self._t_prev = time.perf_counter() - self._e_prev = 0.0 + self._cpu_idle = _IdleEstimator(idle_horizon_s) + #: J per CPU-second charged in the last window, and where it came from + #: (``"model"``, ``"fit"`` or ``"average"``). + self.cpu_j_per_cpu_s: float | None = None + self.cpu_cost_source: str | None = None + self._gpu_idle = _IdleEstimator(idle_horizon_s) + self._prev: WindowSample | None = None + self._t_prev = clock() + self._times_prev = (0.0, 0.0, 0.0) + + def reset_window(self, sample: WindowSample) -> None: + """Anchor the first window at ``sample``. Call when the tracker starts.""" + with self._lock: + self._prev = sample + self._t_prev = sample.timestamp + self._times_prev = self._cpu_times() - 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, meter: _Meter | None = None) -> _InFlight: + """Start weighting a request. Returns the handle to pass to :meth:`end`. - 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()) + ``meter`` is the request's CPU meter; without one it gets no CPU energy. + """ + state = _InFlight(endpoint=endpoint, start=self._clock(), meter=meter) with self._lock: self._in_flight[id(state)] = state return state @@ -100,7 +367,7 @@ def end(self, state: _InFlight) -> None: drop that energy into a zero-width window and silently lose it. """ with self._lock: - state.end = time.perf_counter() + state.end = self._clock() def close(self) -> None: """Emit every in-flight request as-is. Call after the tracker stops.""" @@ -110,15 +377,15 @@ def close(self) -> None: 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. + def on_window(self, sample: WindowSample) -> None: + """Close a sampling window with the tracker's cumulative energies. 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) + self._settle(sample) finished = [ self._in_flight.pop(key) for key, state in list(self._in_flight.items()) @@ -128,58 +395,151 @@ def on_window(self, total_energy_kwh: float) -> None: for state in finished: self._emit(state) - def _settle(self, total_energy_kwh: float) -> None: + def _settle(self, sample: WindowSample) -> None: """Split one window. Caller must hold ``self._lock``.""" - now = time.perf_counter() - w0, w1 = self._t_prev, now + prev, w0, times_prev = self._prev, self._t_prev, self._times_prev + times = self._cpu_times() + w1 = sample.timestamp width = w1 - w0 - delta = total_energy_kwh - self._e_prev - if width <= 0: - self._t_prev, self._e_prev = now, total_energy_kwh + if prev is None: + self._prev, self._t_prev, self._times_prev = sample, w1, times 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 + delta = sample.total_kwh - prev.total_kwh + d_cpu = sample.cpu_kwh - prev.cpu_kwh + d_gpu = sample.gpu_kwh - prev.gpu_kwh + if width <= 0 or min(delta, d_cpu, d_gpu) < 0: + # Counter wraparound or reset, or a window with no width: no + # honest way to split it. The CPU time metered in it is dropped + # with it, or the next window would charge it against a process + # time that excludes it. + # A zero-width window with no energy lost nothing: not counted. + self.windows_skipped += int(width > 0 or delta != 0) + for state in self._in_flight.values(): + state.cpu_seen_ns = state.meter.total_ns() if state.meter else 0 + self._prev, self._t_prev, self._times_prev = sample, w1, times return - states: list[_InFlight] = [] - weights: list[float] = [] - for state in self._in_flight.values(): + d_proc, d_busy, d_total = (now - then for now, then in zip(times, times_prev)) + busy_cpus = d_busy / width if d_total > 0 else None + cpu_idle_w = sample.cpu_idle_w + if cpu_idle_w is None: + cpu_idle_w = self._cpu_idle.update( + w1, d_cpu * _WS_PER_KWH / width, busy_cpus + ) + else: + self._cpu_idle.watts = cpu_idle_w + dynamic_cpu = _dynamic(d_cpu, cpu_idle_w, width) + # This process's share of the machine's busy CPU time. psutil counts in + # clock ticks (10 ms on Linux), so short windows are noisy: clamped. + if sample.cpu_per_process: + share = 1.0 # load mode in process tracking: already ours alone + elif d_busy > 0: + share = min(max(d_proc / d_busy, 0.0), 1.0) + else: + share = 1.0 if d_proc > 0 else 0.0 + ours = dynamic_cpu * share + dynamic_gpu = 0.0 + if sample.gpu_quality is not None: + gpu_idle_w = self._gpu_idle.update(w1, d_gpu * _WS_PER_KWH / width) + dynamic_gpu = _dynamic(d_gpu, gpu_idle_w, width) + + quality = sample.cpu_quality + if sample.gpu_quality is not None: + quality = min(quality, sample.gpu_quality, key=_QUALITIES.index) + + # Cost of one CPU-second, J: the power model's, else the slope of a + # good fit of power on load, else our average (dynamic / CPU time). + # ponytail: one linear slope. On a convex power curve it undercharges + # at high load and overcharges at low load; the cap below bounds the + # second case by what our process actually drew. + d_proc_s = max(d_proc, 0.0) + if sample.cpu_w_per_busy_cpu is not None: + j_per_cpu_s, self.cpu_cost_source = sample.cpu_w_per_busy_cpu, "model" + elif self._cpu_idle.slope is not None: + j_per_cpu_s, self.cpu_cost_source = self._cpu_idle.slope, "fit" + else: + j_per_cpu_s = ours * _WS_PER_KWH / d_proc_s if d_proc_s > 0 else 0.0 + self.cpu_cost_source = "average" + self.cpu_j_per_cpu_s = j_per_cpu_s + + states = list(self._in_flight.values()) + # CPU: each request is charged its metered CPU time at that cost. If + # the meters claim more than the process used (clock granularity), + # they are scaled down to fit. If all of our process's CPU time at that + # cost exceeds our share of the dynamic energy, one cap scales the + # requests and the unclaimed rest down alike. + seen = [state.meter.total_ns() if state.meter else 0 for state in states] + cpu_s = [(now - state.cpu_seen_ns) / 1e9 for now, state in zip(seen, states)] + claimed_s = sum(cpu_s) + if claimed_s > d_proc_s: + cpu_s = [c * d_proc_s / claimed_s for c in cpu_s] + claimed_s = d_proc_s + process_wanted = j_per_cpu_s * d_proc_s / _WS_PER_KWH + cap = min(1.0, ours / process_wanted) if process_wanted > 0 else 1.0 + cpu_parts = [j_per_cpu_s * c / _WS_PER_KWH * cap for c in cpu_s] + cpu_attributed = sum(cpu_parts) + # The rest of our process's CPU time, at the same cost and cap. + process_rest = max(min(process_wanted * cap, ours) - cpu_attributed, 0.0) + other = dynamic_cpu - cpu_attributed - process_rest + # GPU: no per-request signal, so split by overlap with the window. + weights = [] + for state in states: 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) + weights.append(max(hi - lo, 0.0)) + total_weight = sum(weights) + gpu_parts = [ + dynamic_gpu * weight / total_weight if total_weight > 0 else 0.0 + for weight in weights + ] if not states: - self.unattributed_kwh += delta + self.unattributed_kwh += process_rest + dynamic_gpu else: - total_weight = sum(weights) - for state, weight in zip(states, weights): - share = delta * (weight / total_weight) - state.energy += share + gpu_attributed = sum(gpu_parts) + for state, now, cpu_part, gpu_part in zip( + states, seen, cpu_parts, gpu_parts + ): + state.cpu_seen_ns = now + state.energy += cpu_part + state.gpu_energy += gpu_part + state.quality = min( + quality, state.quality or quality, key=_QUALITIES.index + ) state.windows += 1 state.concurrency_sum += len(states) - self.attributed_kwh += share - + self.attributed_kwh += cpu_attributed + gpu_attributed + self.process_unattributed_kwh += process_rest + # No overlap at all (or rounding): nobody to give it to. + self.unattributed_kwh += dynamic_gpu - gpu_attributed + self.other_processes_kwh += other + self.idle_kwh += delta - dynamic_cpu - dynamic_gpu # 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. + # settled_kwh; left in place, the next window covers it. self.windows_settled += 1 self.settled_kwh += delta - self._t_prev, self._e_prev = now, total_energy_kwh + self._prev, self._t_prev, self._times_prev = sample, w1, times def _emit(self, state: _InFlight) -> None: + if state.energy > 0 and state.gpu_energy > 0: + method = "mixed" + elif state.gpu_energy > 0: + method = "wall" + else: + method = "cpu_time" result = RequestEnergy( endpoint=state.endpoint, energy_kwh=state.energy if state.windows else None, - duration_s=(state.end or time.perf_counter()) - state.start, + gpu_kwh=state.gpu_energy if state.windows else None, + duration_s=(state.end or self._clock()) - state.start, windows=state.windows, mean_concurrency=( state.concurrency_sum / state.windows if state.windows else None ), + cpu_seconds=state.meter.total_ns() / 1e9 if state.meter else 0.0, + attribution_method=method, + quality=state.quality or "none", ) if state.on_resolved is not None: try: @@ -191,8 +551,18 @@ def report(self) -> dict[str, Any]: """Run-level accounting, for checking what the split did.""" return { "attributed_kwh": self.attributed_kwh, + "idle_kwh": self.idle_kwh, + "other_processes_kwh": self.other_processes_kwh, + "process_unattributed_kwh": self.process_unattributed_kwh, "unattributed_kwh": self.unattributed_kwh, "settled_kwh": self.settled_kwh, + "idle_power_w": {"cpu": self._cpu_idle.watts, "gpu": self._gpu_idle.watts}, + "cpu_j_per_cpu_second": self.cpu_j_per_cpu_s, + "cpu_cost_source": self.cpu_cost_source, + "quality": { + "cpu": self._prev.cpu_quality if self._prev else None, + "gpu": self._prev.gpu_quality if self._prev else None, + }, "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 index e44ccc527..5d37b9e4e 100644 --- a/codecarbon/integrations/fastapi/middleware.py +++ b/codecarbon/integrations/fastapi/middleware.py @@ -3,9 +3,13 @@ from __future__ import annotations import functools +import threading from collections.abc import Callable +from contextvars import ContextVar +from typing import Any try: + import anyio.to_thread from starlette.types import ASGIApp, Message, Receive, Scope, Send except ImportError as e: # pragma: no cover raise ImportError( @@ -13,9 +17,68 @@ "pip install 'codecarbon[fastapi]'" ) from e -from codecarbon.emissions_tracker import BaseEmissionsTracker +from codecarbon.emissions_tracker import BaseEmissionsTracker, WindowSample from codecarbon.external.logger import logger -from codecarbon.integrations.fastapi.attribution import EnergyAttributor, RequestEnergy +from codecarbon.integrations.fastapi.attribution import ( + EnergyAttributor, + RequestEnergy, + _Meter, + _Metered, +) + +#: The meter of the request whose task (or child task) is running, if any. +_current_meter: ContextVar[_Meter | None] = ContextVar( + "codecarbon_request_meter", default=None +) +_patch_lock = threading.Lock() +_patch_users = 0 +#: Whether ``_metered_run_sync`` is reachable from ``anyio.to_thread.run_sync``. +_patch_in_chain = False +_original_run_sync: Callable[..., Any] = anyio.to_thread.run_sync + + +async def _metered_run_sync(func: Callable[..., Any], *args: Any, **kwargs: Any): + """``anyio.to_thread.run_sync`` that meters the call for the current request. + + Starlette and FastAPI send sync endpoints, sync dependencies, sync + iterators and file work through it. Outside a metered request it calls + the original unchanged. + """ + meter = _current_meter.get() + if meter is not None: + func = functools.partial(meter.run_in_thread, func) + return await _original_run_sync(func, *args, **kwargs) + + +def _patch_run_sync() -> None: + """Route ``anyio.to_thread.run_sync`` through the meter. Reference-counted. + + Once patched, ours is assumed to stay in the chain. If a library that + wrapped ``run_sync`` before us later restores its own saved original, + our wrapper is dropped and threadpool work goes unmetered; there is no + cheap way to tell whether a wrapper on top still calls ours. + """ + global _patch_users, _patch_in_chain, _original_run_sync + with _patch_lock: + if not _patch_in_chain: + _original_run_sync = anyio.to_thread.run_sync + anyio.to_thread.run_sync = _metered_run_sync + _patch_in_chain = True + _patch_users += 1 + + +def _unpatch_run_sync() -> None: + """Undo :func:`_patch_run_sync` once its last user is gone. + + If someone patched ``run_sync`` on top of ours, ours stays in their chain + (it does nothing outside a metered request) rather than breaking it. + """ + global _patch_users, _patch_in_chain + with _patch_lock: + _patch_users -= 1 + if _patch_users == 0 and anyio.to_thread.run_sync is _metered_run_sync: + anyio.to_thread.run_sync = _original_run_sync + _patch_in_chain = False def log_request( @@ -23,9 +86,10 @@ def log_request( ) -> None: """Default ``on_request`` handler; logs via the ``codecarbon`` logger.""" logger.debug( - "CodeCarbon %s: energy=%s kWh emissions=%s kg CO2 status=%s", + "CodeCarbon %s: cpu=%s kWh gpu=%s kWh emissions=%s kg CO2 status=%s", energy.endpoint, energy.energy_kwh, + energy.gpu_kwh, emissions_kg, status_code, ) @@ -48,6 +112,13 @@ class CodeCarbonMiddleware: tracker: Optional tracker; defaults to ``app.state.codecarbon_tracker``. on_request: Callback ``(RequestEnergy, emissions_kg | None, status_code)``. ``None`` disables reporting. + meter_threadpool: Also meter the CPU time of work a request sends to + the threadpool (sync endpoints and dependencies, sync iterators). + This wraps ``anyio.to_thread.run_sync`` process-wide while a + tracker is attached; outside a request it behaves as before. + There is one current meter per request context, so with nested + ``CodeCarbonMiddleware`` instances the innermost one takes the + threadpool CPU time and the outer ones lose it. """ def __init__( @@ -58,12 +129,16 @@ def __init__( on_request: ( Callable[[RequestEnergy, float | None, int], None] | None ) = log_request, + meter_threadpool: bool = True, ) -> None: self.app = app + self.meter_threadpool = meter_threadpool + self._patched = False self.tracker = tracker self.on_request = on_request self.attributor = EnergyAttributor() self._attached: BaseEmissionsTracker | None = None + self._observer: Callable[[WindowSample], None] | None = None # kg CO2eq per kWh, refreshed once per sampling window. self._intensity: float | None = None @@ -73,17 +148,24 @@ def close(self) -> None: Called automatically on lifespan shutdown. """ if self._attached is not None: - self._attached.remove_energy_window_observer(self._on_window) - self._attached = None + self._attached.remove_energy_window_observer(self._observer) + self._attached = self._observer = None + if self._patched: + _unpatch_run_sync() + self._patched = False self.attributor.close() - def _on_window(self, total_energy_kwh: float) -> None: + def _on_window(self, tracker: BaseEmissionsTracker, sample: WindowSample) -> None: + # A detached tracker's scheduler may already be inside this call: its + # sample must not mix into the next tracker's anchor. + if tracker is not self._attached: + return # Scheduler thread: one intensity lookup per window, not per request. try: - self._intensity = self._attached._carbon_intensity_kg_per_kwh() + self._intensity = tracker._carbon_intensity_kg_per_kwh() except Exception: logger.debug("CodeCarbon: carbon intensity unavailable", exc_info=True) - self.attributor.on_window(total_energy_kwh) + self.attributor.on_window(sample) def _running_tracker(self, scope: Scope) -> BaseEmissionsTracker | None: tracker = self.tracker @@ -114,9 +196,13 @@ async def lifespan_send(message: Message) -> None: # 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.attributor.reset_window(tracker._window_sample()) self._attached = tracker + self._observer = functools.partial(self._on_window, tracker) + tracker.add_energy_window_observer(self._observer) + if self.meter_threadpool: + _patch_run_sync() + self._patched = True if tracker is None: await self.app(scope, receive, send) return @@ -129,10 +215,14 @@ async def send_wrapper(message: Message) -> None: status_code = message["status"] await send(message) - state = self.attributor.begin(_endpoint(scope)) + meter = _Meter() + state = self.attributor.begin(_endpoint(scope), meter) + token = _current_meter.set(meter) if self._patched else None try: - await self.app(scope, receive, send_wrapper) + await _Metered(self.app(scope, receive, send_wrapper), meter) finally: + if token is not None: + _current_meter.reset(token) # 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) @@ -143,7 +233,7 @@ def _resolved(self, status_code: int, energy: RequestEnergy) -> None: if self.on_request is None: return emissions_kg = ( - energy.energy_kwh * self._intensity + (energy.energy_kwh + energy.gpu_kwh) * self._intensity if energy.energy_kwh is not None and self._intensity is not None else None ) diff --git a/docs/how-to/examples.md b/docs/how-to/examples.md index 60b2a2122..7197805a0 100644 --- a/docs/how-to/examples.md +++ b/docs/how-to/examples.md @@ -160,11 +160,11 @@ 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. +One tracker runs for the app lifetime; the middleware splits the energy of +each of its sampling windows across the requests that were in flight during +that window. 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), @@ -201,10 +201,106 @@ tracker's scheduler thread — keep it cheap. The default callback logs at DEBUG 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. +Each window is split as follows: + +1. Per component, idle power times the window length is taken out first. + Three CPU modes fix it: CPU load mode tracking the machine uses 10% of + TDP, CPU load mode tracking the process uses 0 W, and constant (TDP) mode + counts all of its power as idle, so it charges requests nothing. Otherwise + idle power is estimated as the lower of the lowest window power seen over + the last hour and, when a power-against-CPU-utilisation fit is good + (R² > 0.8, at least 20 windows), that fit's intercept. All RAM energy + counts as idle. +2. Of the CPU energy above idle, this process keeps its share of the + machine's busy CPU time (`psutil.cpu_times()`). The rest belongs to other + processes on the host. +3. Each request is charged the CPU time it used, measured by + `time.thread_time_ns()` around every step of the request's coroutine and + around each call it sends to the threadpool, times the cost of one + CPU-second. That cost comes from the power model in CPU load mode, from + the slope of the power-against-load fit when it is good, and otherwise + from the process's average (its energy above idle over its CPU time). + The process's CPU time that no request claimed goes to + `process_unattributed_kwh` at the same cost. If all of the process's CPU + time at that cost would exceed its share from step 2, requests and that + unclaimed rest are scaled down by the same factor. Whatever the share + leaves over is counted under other processes. +4. GPU energy above idle has no per-request signal, so it is split by how long + each request overlapped the window. + +The fields of `RequestEnergy`: + +| Field | Meaning | +|---|---| +| `energy_kwh` | CPU energy above idle, charged by the request's CPU time | +| `gpu_kwh` | GPU energy above idle, charged by wall-clock overlap | +| `cpu_seconds` | CPU time the meter saw the request use | +| `attribution_method` | `cpu_time`, `wall` (GPU energy only) or `mixed` | +| `quality` | `measured` (RAPL, powermetrics, NVML), `modeled` (CPU load mode) or `none` (constant TDP), for the weakest component used | +| `windows`, `mean_concurrency` | Windows the request was in flight for, and how many requests shared them | + +`emissions_kg` covers `energy_kwh + gpu_kwh`. Both are `None` only when the +request never covered a completed sampling window, which in practice means it +was still pending when the tracker stopped. + +The middleware's `attributor.report()` returns the run-level buckets. They add +up to `settled_kwh` (up to float rounding): + +- `attributed_kwh`: charged to requests. +- `idle_kwh`: idle power and RAM. +- `other_processes_kwh`: CPU energy of other processes on the host. +- `process_unattributed_kwh`: this process's CPU energy that no request meter + claimed. A large value means much of the work runs where the meter cannot + see it (see the limits below). +- `unattributed_kwh`: in windows with no request in flight, this process's + CPU energy above idle (at the cost of a CPU-second) and all GPU energy + above idle; in other windows, GPU energy above idle that no request + overlapped. + +It also reports the idle power in use per component (`idle_power_w`), the +cost of a CPU-second and its source (`cpu_j_per_cpu_second`, +`cpu_cost_source`), and the source quality of the last window. + +Limits: + +- One CPU-second is priced with a single linear slope. Power curves are + convex, so this overcharges at low machine load (the charges are then + capped at the process's share) and undercharges at high load. The upside + is that a request's energy barely moves when another process loads the + host. In CPU load mode the slope is the chord of the cubic model, so the + cap nearly always binds at low load and charges fall back to the average + cost. There, a request's energy rises when other processes load the + machine (by several times in our tests). Only a measured source (RAPL, + powermetrics, NVML, `quality == "measured"`) gives numbers that hold up + under a busy neighbour; the accuracy tests only run against one. +- CPU time is a proxy for energy. Frequency scaling, SMT and wide vector + instructions make one CPU second cost different amounts of energy; the + error from this has not been measured yet. +- The request's own coroutine is metered, and so is the work it sends to the + threadpool: sync (`def`) endpoints and dependencies, sync iterators of a + `StreamingResponse`. For the latter the middleware wraps + `anyio.to_thread.run_sync` process-wide while a tracker is attached; it + changes nothing outside a request. Pass `meter_threadpool=False` to turn + this off. On Linux a long sync call is charged window by window; on other + platforms its CPU time counts in the window where it returns. +- The threadpool wrapper keeps one current meter per request context. With + nested `CodeCarbonMiddleware` instances the innermost one takes the + threadpool CPU time and the outer ones lose it. If another library wrapped + `anyio.to_thread.run_sync` before the middleware did and later restores its + own saved original, the middleware's wrapper drops out of the chain and + threadpool work goes unmetered until the process restarts. +- In CPU load mode with `tracking_mode="process"`, the tracker's energy + includes child processes, but the process CPU time used for pricing does + not. The children's CPU energy lands in `other_processes_kwh`. +- Tasks the request starts with `asyncio.create_task` or `gather` are not + metered (their sync threadpool calls are). Their CPU time, and threads the + application starts itself, land in `process_unattributed_kwh`. +- GPU energy is split by wall-clock time, not by the work each request sent + to the GPU. +- The idle estimate needs time to settle. Until the host has had a quiet + window, idle power is overestimated and requests are undercharged. +- `psutil.cpu_times()` counts in clock ticks (10 ms on Linux), so the process + share is noisy for short windows. With RAPL, a `measure_power_secs` between + 1 and 5 is a good trade-off between window resolution and that noise. + +Sum the values per route over many requests rather than reading a single one. diff --git a/pyproject.toml b/pyproject.toml index 8471ec54f..cf553d6a6 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -165,6 +165,9 @@ version_pattern = "MAJOR.MINOR.PATCH[_TAGNUM]" [tool.pytest.ini_options] pythonpath = "." +markers = [ + "accuracy: bare-metal accuracy checks, opt-in (needs Linux with readable RAPL)", +] [tool.coverage.run] source = [ diff --git a/tests/integrations/fastapi_accuracy.py b/tests/integrations/fastapi_accuracy.py new file mode 100644 index 000000000..eeead260d --- /dev/null +++ b/tests/integrations/fastapi_accuracy.py @@ -0,0 +1,145 @@ +"""Bare-metal accuracy checks for the FastAPI per-request attribution. + +Opt-in, and only meaningful on Linux with readable RAPL counters; skips +otherwise. Takes about a minute and a half:: + + uv run pytest tests/integrations/fastapi_accuracy.py -m accuracy -s + +The app and the load generator share one process and one event loop, so no +server dependency is needed. The client's CPU time is outside every request's +meter: it lands in ``process_unattributed_kwh`` and does not bias the +per-request numbers. The CPU hog runs in separate processes. +""" + +import asyncio +import multiprocessing +import statistics +import time +from collections import defaultdict + +import pytest + +from codecarbon.core.cpu import is_rapl_available +from codecarbon.core.util import is_linux_os + +pytestmark = [ + pytest.mark.accuracy, + pytest.mark.skipif( + not (is_linux_os() and is_rapl_available()), reason="needs readable RAPL" + ), +] + +CONCURRENCY = 16 +PHASE_S = 30.0 +ROUTES = ("/burn1", "/burn2", "/burn4", "/sleep4") + + +def _burn(seconds: float) -> None: + end = time.thread_time_ns() + int(seconds * 1e9) + while time.thread_time_ns() < end: + pass + + +def _hog(stop) -> None: + while not stop.is_set(): + _burn(0.01) + + +def _app(): + from fastapi import FastAPI + + from codecarbon.integrations.fastapi import CodeCarbonMiddleware + + seen = defaultdict(list) + app = FastAPI() + + def record(energy, kg, status): + if energy.energy_kwh is not None: + seen[energy.endpoint.split()[-1]].append(energy.energy_kwh) + + app.add_middleware(CodeCarbonMiddleware, on_request=record) + + for ms in (1, 2, 4): + + async def burn(ms=ms): + _burn(ms / 1000) + return {} + + app.get(f"/burn{ms}")(burn) + + @app.get("/sleep4") + async def sleep(): + await asyncio.sleep(0.004) + return {} + + return app, seen + + +async def _load(app, seconds: float) -> None: + import httpx + + transport = httpx.ASGITransport(app=app) + async with httpx.AsyncClient(transport=transport, base_url="http://t") as client: + deadline = time.monotonic() + seconds + + async def worker(n: int) -> None: + while time.monotonic() < deadline: + assert (await client.get(ROUTES[n % len(ROUTES)])).status_code == 200 + n += 1 + + await asyncio.gather(*(worker(i) for i in range(CONCURRENCY))) + + +def _phase(hog_processes: int = 0) -> dict[str, float]: + """Mean per-request energy per route, kWh.""" + from codecarbon.emissions_tracker import OfflineEmissionsTracker + + tracker = OfflineEmissionsTracker( + country_iso_code="FRA", + measure_power_secs=1, + output_methods=[], + allow_multiple_runs=True, + ) + tracker.start() + if tracker._window_sample().cpu_quality != "measured": + tracker.stop() + pytest.skip("the tracker did not pick a measured CPU source") + stop = multiprocessing.Event() + hogs = [ + multiprocessing.Process(target=_hog, args=(stop,)) for _ in range(hog_processes) + ] + for hog in hogs: + hog.start() + app, seen = _app() + app.state.codecarbon_tracker = tracker + try: + asyncio.run(_load(app, PHASE_S)) + time.sleep(2.5) # let the last requests resolve on real windows + finally: + stop.set() + for hog in hogs: + hog.join() + tracker.stop() + means = {route: statistics.fmean(values) for route, values in seen.items()} + print(f"\nhogs={hog_processes}", {k: f"{v:.3e} kWh" for k, v in means.items()}) + return means + + +@pytest.fixture(scope="module") +def quiet(): + return _phase() + + +def test_energy_follows_cpu_time_1_2_4(quiet): + one = quiet["/burn1"] + assert quiet["/burn2"] / one == pytest.approx(2, rel=0.10) + assert quiet["/burn4"] / one == pytest.approx(4, rel=0.10) + + +def test_sleep_costs_at_most_2_percent_of_a_4ms_burn(quiet): + assert quiet["/sleep4"] <= 0.02 * quiet["/burn4"] + + +def test_cpu_hog_changes_per_request_energy_at_most_15_percent(quiet): + hogged = _phase(hog_processes=max(multiprocessing.cpu_count() // 2, 1)) + assert hogged["/burn4"] == pytest.approx(quiet["/burn4"], rel=0.15) diff --git a/tests/integrations/fastapi_overhead.py b/tests/integrations/fastapi_overhead.py new file mode 100644 index 000000000..0d7184d52 --- /dev/null +++ b/tests/integrations/fastapi_overhead.py @@ -0,0 +1,83 @@ +"""Per-request overhead of ``CodeCarbonMiddleware``. + +Times 100k in-process ASGI calls of a trivial app, with and without the +middleware, and prints p50/p99 per call. Budget: at most 30 µs (Linux) or +50 µs (macOS) added at p50. Run with:: + + uv run python tests/integrations/fastapi_overhead.py + +A sampling window is closed every 1000 requests, outside the timed calls, so +the in-flight table stays at a realistic size. +""" + +import asyncio +import statistics +import sys +import time + +from codecarbon.emissions_tracker import OfflineEmissionsTracker +from codecarbon.integrations.fastapi import CodeCarbonMiddleware + +N = 100_000 +WINDOW_EVERY = 1000 +SCOPE = {"type": "http", "method": "GET", "path": "/", "headers": []} + + +async def _app(scope, receive, send): + await send({"type": "http.response.start", "status": 200, "headers": []}) + await send({"type": "http.response.body", "body": b"ok"}) + + +async def _receive(): + return {"type": "http.request", "body": b"", "more_body": False} + + +async def _send(message): + pass + + +async def _time(app, on_window=None) -> list[int]: + durations = [] + for i in range(N): + start = time.perf_counter_ns() + await app(dict(SCOPE), _receive, _send) + durations.append(time.perf_counter_ns() - start) + if on_window is not None and i % WINDOW_EVERY == WINDOW_EVERY - 1: + on_window() + return durations + + +def _summary(name: str, durations: list[int]) -> float: + q = statistics.quantiles(durations, n=100) + p50, p99 = q[49] / 1000, q[98] / 1000 + print(f"{name:>16}: p50 {p50:7.2f} µs p99 {p99:7.2f} µs") + return p50 + + +def main() -> None: + tracker = OfflineEmissionsTracker( + country_iso_code="FRA", + measure_power_secs=3600, # windows are closed by hand + output_methods=[], + allow_multiple_runs=True, + log_level="error", + ) + tracker.start() + try: + middleware = CodeCarbonMiddleware(_app, tracker=tracker, on_request=None) + asyncio.run(_time(_app)) # warm-up + bare = _summary("bare", asyncio.run(_time(_app))) + metered = _summary( + "with middleware", + asyncio.run(_time(middleware, tracker._measure_power_and_energy)), + ) + finally: + tracker.stop() + budget = 30.0 if sys.platform.startswith("linux") else 50.0 + added = metered - bare + print(f"{'added':>16}: p50 {added:7.2f} µs (budget {budget:.0f} µs)") + sys.exit(0 if added <= budget else 1) + + +if __name__ == "__main__": + main() diff --git a/tests/integrations/test_fastapi.py b/tests/integrations/test_fastapi.py index 403949efe..e67924ed5 100644 --- a/tests/integrations/test_fastapi.py +++ b/tests/integrations/test_fastapi.py @@ -1,67 +1,196 @@ """Tests for the FastAPI per-request energy attribution.""" +import asyncio import threading import time +import types from contextlib import asynccontextmanager from types import SimpleNamespace +import psutil import pytest from fastapi import FastAPI from fastapi.testclient import TestClient -from codecarbon.emissions_tracker import OfflineEmissionsTracker +from codecarbon.emissions_tracker import OfflineEmissionsTracker, WindowSample +from codecarbon.external.hardware import CPU from codecarbon.integrations.fastapi import ( CodeCarbonMiddleware, EnergyAttributor, RequestEnergy, ) +from codecarbon.integrations.fastapi.attribution import _Meter, _Metered + + +def _sample(total: float, *, cpu: float | None = None, gpu: float = 0.0, **kw): + """A window sample; by default all energy is dynamic CPU energy.""" + fields = dict( + timestamp=time.perf_counter(), + total_kwh=total, + cpu_kwh=total - gpu if cpu is None else cpu, + gpu_kwh=gpu, + ram_kwh=0.0, + cpu_quality="measured", + gpu_quality=None, + cpu_idle_w=0.0, + cpu_w_per_busy_cpu=None, + cpu_per_process=False, + ) + fields.update(kw) + return WindowSample(**fields) + + +def _all_ours(): + """CPU times where every busy second is this process's (share 1).""" + n = 0.0 + + def times(): + nonlocal n + n += 1.0 + return n, n, 2 * n + + return times def _invariant(attributor: EnergyAttributor) -> None: report = attributor.report() - assert report["attributed_kwh"] + report["unattributed_kwh"] == pytest.approx( + buckets = ( + "attributed_kwh", + "idle_kwh", + "other_processes_kwh", + "process_unattributed_kwh", + "unattributed_kwh", + ) + assert sum(report[b] for b in buckets) == 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) + attributor = EnergyAttributor(cpu_times=_all_ours()) + attributor.reset_window(_sample(0.0)) + attributor.on_window(_sample(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) +def _begin(attributor, endpoint="GET /a", *, start=0.0, cpu_s=0.0): + """A request that started at ``start`` and used ``cpu_s`` of CPU so far.""" + meter = _Meter() + meter.ns = int(cpu_s * 1e9) + state = attributor.begin(endpoint, meter) + state.start = start + return state + + +def test_cpu_energy_splits_by_cpu_time(): + times = _FakeTimes() + attributor = EnergyAttributor(cpu_times=times) + attributor.reset_window(_sample(0.0, timestamp=0.0)) + busy = _begin(attributor, "GET /busy", cpu_s=0.6) + waiting = _begin(attributor, "GET /io", cpu_s=0.0) + times.advance(process=1.0, busy=1.0) + attributor.on_window(_sample(_kwh(100, 1), timestamp=1.0)) + # The process used 1 s of CPU; the meters account for 0.6 s of it. + assert busy.energy == pytest.approx(_kwh(60, 1)) + assert waiting.energy == 0.0 + assert attributor.process_unattributed_kwh == pytest.approx(_kwh(40, 1)) + # Only the CPU time since the last window counts in the next one. + busy.meter.ns += int(0.5e9) + times.advance(process=0.5, busy=0.5) + attributor.on_window(_sample(_kwh(150, 1), timestamp=2.0)) + assert busy.energy == pytest.approx(_kwh(110, 1)) + _invariant(attributor) + + +def test_meters_claiming_more_than_the_process_are_scaled_down(): + times = _FakeTimes() + attributor = EnergyAttributor(cpu_times=times) + attributor.reset_window(_sample(0.0, timestamp=0.0)) + a = _begin(attributor, cpu_s=0.9) + b = _begin(attributor, cpu_s=0.3) + times.advance(process=1.0, busy=1.0) + attributor.on_window(_sample(_kwh(100, 1), timestamp=1.0)) + assert a.energy == pytest.approx(_kwh(75, 1)) + assert b.energy == pytest.approx(_kwh(25, 1)) + assert attributor.process_unattributed_kwh == pytest.approx(0.0, abs=1e-18) + _invariant(attributor) + - # `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) +def test_gpu_energy_splits_by_overlap(): + attributor = EnergyAttributor(cpu_times=_all_ours()) + gpu = dict(gpu_quality="measured") + attributor.reset_window(_sample(0.0, timestamp=0.0, **gpu)) + attributor.on_window(_sample(0.0, timestamp=1.0, **gpu)) # GPU idles at 0 W + early = _begin(attributor, "GET /a", start=1.0) + late = _begin(attributor, "GET /b", start=2.0) + attributor.on_window(_sample(_kwh(90, 2), gpu=_kwh(90, 2), timestamp=3.0, **gpu)) + # `early` overlapped twice as much of the window as `late`. + assert early.gpu_energy == pytest.approx(_kwh(120, 1)) + assert late.gpu_energy == pytest.approx(_kwh(60, 1)) + assert early.energy == late.energy == 0.0 # no CPU time metered _invariant(attributor) +def test_request_reports_method_and_weakest_quality(): + times = _FakeTimes() + attributor = EnergyAttributor(cpu_times=times) + gpu = dict(gpu_quality="measured", cpu_quality="modeled") + attributor.reset_window(_sample(0.0, timestamp=0.0, **gpu)) + attributor.on_window(_sample(0.0, timestamp=1.0, **gpu)) + results = [] + both = _begin(attributor, "GET /both", start=1.0, cpu_s=0.5) + gpu_only = _begin(attributor, "GET /gpu", start=1.0) + both.on_resolved = gpu_only.on_resolved = results.append + attributor.end(both) + attributor.end(gpu_only) + times.advance(process=1.0, busy=1.0) + attributor.on_window(_sample(_kwh(20, 1), gpu=_kwh(10, 1), timestamp=2.0, **gpu)) + by_endpoint = {r.endpoint: r for r in results} + assert by_endpoint["GET /both"].attribution_method == "mixed" + assert by_endpoint["GET /both"].energy_kwh == pytest.approx(_kwh(5, 1)) + assert by_endpoint["GET /both"].gpu_kwh == pytest.approx(_kwh(5, 1)) + assert by_endpoint["GET /both"].cpu_seconds == 0.5 + assert by_endpoint["GET /gpu"].attribution_method == "wall" + assert {r.quality for r in results} == {"modeled"} + assert attributor.report()["quality"] == {"cpu": "modeled", "gpu": "measured"} + + 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 + times = _FakeTimes() + attributor = EnergyAttributor(cpu_times=times) + attributor.reset_window(_sample(5.0, timestamp=0.0)) + state = _begin(attributor, cpu_s=1.0) + times.advance(process=1.0, busy=1.0) + attributor.on_window(_sample(1.0, timestamp=1.0)) # RAPL wrap assert attributor.windows_skipped == 1 assert attributor.attributed_kwh == 0.0 + # The skipped window's CPU time is not charged to the next one. + times.advance(process=1.0, busy=1.0) + attributor.on_window(_sample(1.0 + _kwh(100, 1), timestamp=2.0)) + assert state.energy == 0.0 + _invariant(attributor) + + +def test_zero_width_window_is_skipped_with_its_cpu_time(): + times = _FakeTimes() + attributor = EnergyAttributor(cpu_times=times) + attributor.reset_window(_sample(0.0, timestamp=1.0)) + state = _begin(attributor, cpu_s=1.0) + times.advance(process=1.0, busy=1.0) + attributor.on_window(_sample(_kwh(10, 1), timestamp=1.0)) # clock stood still + assert attributor.windows_skipped == 1 + times.advance(process=1.0, busy=1.0) + attributor.on_window(_sample(_kwh(10, 1) + _kwh(100, 1), timestamp=2.0)) + assert state.energy == 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) + attributor = EnergyAttributor(cpu_times=_all_ours()) + attributor.reset_window(_sample(0.0)) results = [] state = attributor.begin("GET /fast") state.on_resolved = results.append @@ -73,8 +202,8 @@ def test_unresolved_request_reports_no_energy(): def test_invariant_holds_under_concurrency(): """The core property: nothing is created or lost by the split.""" - attributor = EnergyAttributor() - attributor.reset_window(0.0) + attributor = EnergyAttributor(cpu_times=_all_ours()) + attributor.reset_window(_sample(0.0)) results: list[RequestEnergy] = [] # list.append is atomic under the GIL stop = threading.Event() energy = 0.0 @@ -82,15 +211,18 @@ def test_invariant_holds_under_concurrency(): def sampler(): nonlocal energy while not stop.is_set(): - energy += 0.001 - attributor.on_window(energy) + energy += 0.001 + 0.0005 * (energy * 1000 % 3) + gpu = energy / 3 + attributor.on_window(_sample(energy, gpu=gpu, gpu_quality="measured")) _invariant(attributor) time.sleep(0.002) def requester(i: int): for _ in range(20): - state = attributor.begin(f"GET /{i % 3}") + meter = _Meter() + state = attributor.begin(f"GET /{i % 3}", meter) state.on_resolved = results.append + meter.ns += 5_000_000 time.sleep(0.001) attributor.end(state) @@ -107,11 +239,318 @@ def requester(i: int): _invariant(attributor) assert len(results) == 8 * 20 - resolved = [r.energy_kwh for r in results if r.energy_kwh is not None] + resolved = [r.energy_kwh + r.gpu_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 _FakeTimes: + """Scripted ``(process, busy, total)`` CPU seconds; advanced by hand.""" + + def __init__(self) -> None: + self.process = self.busy = self.total = 0.0 + + def advance(self, *, process=0.0, busy=0.0, total=1.0) -> None: + self.process += process + self.busy += busy + self.total += total + + def __call__(self): + return self.process, self.busy, self.total + + +def _kwh(watts: float, seconds: float) -> float: + return watts * seconds / 3.6e6 + + +def test_idle_power_is_the_rolling_minimum(): + attributor = EnergyAttributor(cpu_times=_all_ours()) + attributor.reset_window(_sample(0.0, timestamp=0.0, cpu_idle_w=None)) + attributor.on_window(_sample(_kwh(10, 1), timestamp=1.0, cpu_idle_w=None)) + # First window: its own power is the only minimum, so it is all idle. + assert attributor.idle_kwh == pytest.approx(_kwh(10, 1)) + attributor.on_window(_sample(_kwh(10 + 40, 1), timestamp=2.0, cpu_idle_w=None)) + assert attributor.unattributed_kwh == pytest.approx(_kwh(30, 1)) + assert attributor.report()["idle_power_w"]["cpu"] == pytest.approx(10) + _invariant(attributor) + + +def test_idle_power_uses_a_good_regression_intercept(): + """Never idle, so the minimum (13 W) overestimates; the fit finds 5 W.""" + times = _FakeTimes() + attributor = EnergyAttributor(cpu_times=times) + attributor.reset_window(_sample(0.0, timestamp=0.0, cpu_idle_w=None)) + energy = 0.0 + for i in range(1, 31): + util = 0.2 + 0.6 * (i % 5) / 4 + times.advance(busy=util) + energy += _kwh(5 + 40 * util, 1) + attributor.on_window(_sample(energy, timestamp=float(i), cpu_idle_w=None)) + assert attributor.report()["idle_power_w"]["cpu"] == pytest.approx(5) + _invariant(attributor) + + +def test_idle_rolling_minimum_forgets_old_windows(): + attributor = EnergyAttributor(cpu_times=_FakeTimes(), idle_horizon_s=10) + attributor.reset_window(_sample(0.0, timestamp=0.0, cpu_idle_w=None)) + attributor.on_window(_sample(_kwh(10, 1), timestamp=1.0, cpu_idle_w=None)) + energy = _kwh(10, 1) + for t in range(2, 20): + energy += _kwh(20, 1) + attributor.on_window(_sample(energy, timestamp=float(t), cpu_idle_w=None)) + assert attributor.report()["idle_power_w"]["cpu"] == pytest.approx(20) + + +def test_analytic_idle_and_ram_go_to_idle_bucket(): + """Load mode fixes idle at 0.1 * TDP; RAM is never charged to requests.""" + attributor = EnergyAttributor(cpu_times=_all_ours()) + attributor.reset_window(_sample(0.0, timestamp=0.0)) + state = _begin(attributor, cpu_s=1.0) # all of the process's CPU time + cpu, ram = _kwh(25, 2), _kwh(3, 2) + attributor.on_window( + _sample(cpu + ram, cpu=cpu, ram_kwh=ram, timestamp=2.0, cpu_idle_w=10.0) + ) + assert state.energy == pytest.approx(_kwh(15, 2)) + assert attributor.idle_kwh == pytest.approx(_kwh(10 + 3, 2)) + _invariant(attributor) + + +def test_gpu_idle_is_estimated_separately(): + attributor = EnergyAttributor(cpu_times=_all_ours()) + gpu = dict(gpu_quality="measured") + attributor.reset_window(_sample(0.0, timestamp=0.0, **gpu)) + attributor.on_window(_sample(_kwh(50, 1), gpu=_kwh(50, 1), timestamp=1.0, **gpu)) + attributor.on_window( + _sample(_kwh(50 + 100, 1), gpu=_kwh(50 + 100, 1), timestamp=2.0, **gpu) + ) + assert attributor.report()["idle_power_w"]["gpu"] == pytest.approx(50) + assert attributor.unattributed_kwh == pytest.approx(_kwh(50, 1)) + _invariant(attributor) + + +def test_other_processes_keep_their_share_of_cpu_energy(): + """Fake psutil: we used 1 of the machine's 4 busy CPU seconds.""" + times = _FakeTimes() + attributor = EnergyAttributor(cpu_times=times) + attributor.reset_window(_sample(0.0, timestamp=0.0)) + state = _begin(attributor, cpu_s=1.0) + times.advance(process=1.0, busy=4.0, total=8.0) + attributor.on_window(_sample(_kwh(40, 1), timestamp=1.0)) + assert state.energy == pytest.approx(_kwh(10, 1)) + assert attributor.other_processes_kwh == pytest.approx(_kwh(30, 1)) + _invariant(attributor) + + +def test_model_cost_per_cpu_second_is_charged_and_capped(): + """Charged b * cpu_seconds; the rest of the machine's energy is not ours.""" + times = _FakeTimes() + attributor = EnergyAttributor(cpu_times=times) + attributor.reset_window(_sample(0.0, timestamp=0.0, cpu_w_per_busy_cpu=10.0)) + state = _begin(attributor, cpu_s=1.0) + # 100 J above idle, 4 busy CPU-seconds, 1 of them ours: our share is 25 J, + # but the model says one CPU-second costs 10 J. + times.advance(process=1.0, busy=4.0, total=8.0) + attributor.on_window(_sample(_kwh(100, 1), timestamp=1.0, cpu_w_per_busy_cpu=10.0)) + assert state.energy == pytest.approx(_kwh(10, 1)) + assert attributor.other_processes_kwh == pytest.approx(_kwh(90, 1)) + assert attributor.report()["cpu_cost_source"] == "model" + # A cost above our share is capped at the share. + state.meter.ns += int(1e9) + times.advance(process=1.0, busy=4.0, total=8.0) + attributor.on_window(_sample(_kwh(200, 1), timestamp=2.0, cpu_w_per_busy_cpu=60.0)) + assert state.energy == pytest.approx(_kwh(10 + 25, 1)) + _invariant(attributor) + + +def test_cost_cap_is_shared_by_requests_and_process_rest(): + """A binding cap scales metered and unmetered process CPU alike.""" + times = _FakeTimes() + attributor = EnergyAttributor(cpu_times=times) + attributor.reset_window(_sample(0.0, timestamp=0.0, cpu_w_per_busy_cpu=60.0)) + state = _begin(attributor, cpu_s=0.5) + # Our share is 25 J for 1 CPU-second; the model asks 60 J for it, so the + # cap binds, and the meters claimed half of our CPU time. + times.advance(process=1.0, busy=4.0, total=8.0) + attributor.on_window(_sample(_kwh(100, 1), timestamp=1.0, cpu_w_per_busy_cpu=60.0)) + assert state.energy == pytest.approx(_kwh(12.5, 1)) + assert attributor.process_unattributed_kwh == pytest.approx(_kwh(12.5, 1)) + _invariant(attributor) + + +def test_fitted_cost_per_cpu_second_ignores_other_processes_load(): + """Linear machine: 5 W idle + 10 W per busy CPU. A hog doesn't move us.""" + times = _FakeTimes() + attributor = EnergyAttributor(cpu_times=times) + measured = dict(cpu_idle_w=None) + attributor.reset_window(_sample(0.0, timestamp=0.0, **measured)) + energy = 0.0 + for t in range(1, 31): + busy = 1.0 + (t % 4) + times.advance(process=0.5, busy=busy, total=8.0) + energy += _kwh(5 + 10 * busy, 1) + attributor.on_window(_sample(energy, timestamp=float(t), **measured)) + assert attributor.report()["cpu_cost_source"] == "fit" + state = _begin(attributor, start=30.0, cpu_s=0.5) + times.advance(process=0.5, busy=6.0, total=8.0) # a hog takes 5.5 CPUs + attributor.on_window(_sample(energy + _kwh(65, 1), timestamp=31.0, **measured)) + assert state.energy == pytest.approx(_kwh(5, 1)) + assert attributor.report()["cpu_j_per_cpu_second"] == pytest.approx(10) + _invariant(attributor) + + +def test_fitted_cost_is_dropped_without_a_load_reading(): + times = _FakeTimes() + attributor = EnergyAttributor(cpu_times=times) + measured = dict(cpu_idle_w=None) + attributor.reset_window(_sample(0.0, timestamp=0.0, **measured)) + energy = 0.0 + for t in range(1, 31): + busy = 1.0 + (t % 4) + times.advance(process=0.5, busy=busy, total=8.0) + energy += _kwh(5 + 10 * busy, 1) + attributor.on_window(_sample(energy, timestamp=float(t), **measured)) + assert attributor.report()["cpu_cost_source"] == "fit" + times.advance(process=0.5, total=0.0) # no machine CPU times this window + attributor.on_window(_sample(energy + _kwh(20, 1), timestamp=31.0, **measured)) + assert attributor.report()["cpu_cost_source"] == "average" + + +def test_process_share_is_clamped_and_skipped_in_process_mode(): + times = _FakeTimes() + attributor = EnergyAttributor(cpu_times=times) + attributor.reset_window(_sample(0.0, timestamp=0.0)) + # Tick rounding can make our CPU time exceed the machine's busy time. + times.advance(process=1.2, busy=1.0) + attributor.on_window(_sample(_kwh(40, 1), timestamp=1.0)) + assert attributor.other_processes_kwh == 0.0 + # Load mode with tracking_mode="process" already measured only us. + times.advance(process=1.0, busy=4.0) + attributor.on_window(_sample(_kwh(80, 1), timestamp=2.0, cpu_per_process=True)) + assert attributor.other_processes_kwh == 0.0 + assert attributor.unattributed_kwh == pytest.approx(_kwh(80, 1)) + _invariant(attributor) + + +def test_cpu_times_busy_excludes_idle_and_steal(monkeypatch): + from collections import namedtuple + + from codecarbon.integrations.fastapi import attribution + + Times = namedtuple("Times", "user nice system idle iowait steal guest guest_nice") + monkeypatch.setattr( + attribution.psutil, "cpu_times", lambda: Times(5, 1, 2, 10, 3, 4, 2, 1) + ) + _, busy, total = attribution._cpu_times() + assert (busy, total) == (8, 25) + + +def _burn(seconds: float) -> None: + end = time.thread_time_ns() + int(seconds * 1e9) + while time.thread_time_ns() < end: + pass + + +def test_metered_returns_and_counts_only_cpu_time(): + async def work(): + _burn(0.01) + await asyncio.sleep(0.05) + _burn(0.01) + return 42 + + async def main(): + return await _Metered(work(), meter) + + meter = _Meter() + assert asyncio.run(main()) == 42 + assert 0.02 <= meter.ns / 1e9 < 0.03 + + +def test_metered_propagates_errors_and_cancellation(): + meter = _Meter() + + async def boom(): + await asyncio.sleep(0) + raise ValueError("boom") + + async def main(): + await _Metered(boom(), meter) + + with pytest.raises(ValueError): + asyncio.run(main()) + + cleaned = [] + + async def slow(): + try: + await asyncio.sleep(10) + finally: + cleaned.append(True) + + async def cancel_it(): + task = asyncio.ensure_future(_Metered(slow(), meter)) + await asyncio.sleep(0.01) + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + + asyncio.run(cancel_it()) + assert cleaned == [True] + + +def test_metered_passes_thrown_exceptions_the_coroutine_handles(): + async def swallow(): + try: + await asyncio.sleep(10) + except asyncio.CancelledError: + return "handled" + + async def main(): + task = asyncio.ensure_future(_Metered(swallow(), _Meter())) + await asyncio.sleep(0.01) + task.cancel() + return await task + + assert asyncio.run(main()) == "handled" + + +def test_metered_drives_any_awaitable(): + """ASGI only promises an awaitable: a Future or custom ``__await__`` works.""" + + class Custom: + def __await__(self): + yield from asyncio.sleep(0).__await__() + return "custom" + + async def main(): + future = asyncio.get_running_loop().create_future() + asyncio.get_running_loop().call_soon(future.set_result, 7) + return await _Metered(future, _Meter()), await _Metered(Custom(), _Meter()) + + assert asyncio.run(main()) == (7, "custom") + + +def test_metered_close_closes_the_inner_coroutine(): + closed = [] + + @types.coroutine + def pause(): + yield + + async def inner(): + try: + await pause() + finally: + closed.append(True) + + async def outer(): + await _Metered(inner(), _Meter()) + + coro = outer() + coro.send(None) # suspended inside pause() + coro.close() + assert closed == [True] + + class _FakeTracker: """The slice of a tracker the middleware touches; windows closed by hand.""" @@ -129,10 +568,14 @@ def remove_energy_window_observer(self, callback): def _carbon_intensity_kg_per_kwh(self): return 0.5 + def _window_sample(self): + # Per-process energy: the real machine's other processes don't matter. + return _sample(self._total_energy.kWh, cpu_per_process=True) + def window(self, kwh: float) -> None: self._total_energy.kWh += kwh for callback in tuple(self.observers): - callback(self._total_energy.kWh) + callback(self._window_sample()) def _app(tracker, seen, *, lifespan=None): @@ -163,8 +606,10 @@ def test_request_resolves_on_next_window(): 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) + # A sync endpoint's work runs in a worker thread the meter doesn't see. + assert 0.0 <= energy.energy_kwh < 1.0 + assert energy.cpu_seconds > 0 + assert kg == pytest.approx(energy.energy_kwh * 0.5) def test_app_raising_is_reported_as_500(): @@ -184,6 +629,29 @@ def test_close_on_shutdown_emits_pending(): assert tracker.observers == [] +def test_late_window_from_a_replaced_tracker_is_ignored(): + async def ok(scope, receive, send): + await send({"type": "http.response.start", "status": 200, "headers": []}) + await send({"type": "http.response.body", "body": b""}) + + async def noop(message): + pass + + old, new = _FakeTracker(), _FakeTracker() + middleware = CodeCarbonMiddleware(ok, tracker=old, on_request=None) + scope = {"type": "http", "method": "GET", "path": "/"} + asyncio.run(middleware(scope, None, noop)) + (late,) = old.observers + middleware.tracker = new + asyncio.run(middleware(scope, None, noop)) + anchor = middleware.attributor._prev + old._total_energy.kWh = 5.0 + late(old._window_sample()) # the scheduler was already calling it + middleware.close() + assert middleware.attributor._prev is anchor + assert middleware.attributor.settled_kwh == 0.0 + + def test_tracker_not_started_records_nothing(): tracker, seen = _FakeTracker(started=False), [] with TestClient(_app(tracker, seen)) as client: @@ -219,3 +687,333 @@ async def lifespan(app: FastAPI): 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 + + +def test_window_sample_reports_components_and_quality(): + tracker = OfflineEmissionsTracker( + country_iso_code="FRA", + measure_power_secs=3600, + output_methods=[], + allow_multiple_runs=True, + force_mode_cpu_load=True, + ) + samples = [] + tracker.add_energy_window_observer(samples.append) + tracker.start() + try: + tracker._measure_power_and_energy() + finally: + tracker.stop() + sample = samples[0] + assert sample.cpu_quality == "modeled" + (cpu,) = [h for h in tracker._hardware if isinstance(h, CPU)] + assert sample.cpu_idle_w == pytest.approx(0.1 * cpu._tdp * tracker._pue) + assert sample.cpu_w_per_busy_cpu == pytest.approx( + 0.9 * cpu._tdp * tracker._pue / psutil.cpu_count() + ) + parts = sample.cpu_kwh + sample.gpu_kwh + sample.ram_kwh + assert sample.total_kwh == pytest.approx(parts) + assert sample.total_kwh > 0 + + +def test_window_sample_error_never_breaks_measurement(monkeypatch): + tracker = OfflineEmissionsTracker( + country_iso_code="FRA", + measure_power_secs=3600, + output_methods=[], + allow_multiple_runs=True, + force_mode_cpu_load=True, + ) + samples = [] + tracker.add_energy_window_observer(samples.append) + + def broken(): + raise RuntimeError("sample") + + monkeypatch.setattr(tracker, "_window_sample", broken) + tracker.start() + try: + tracker._measure_power_and_energy() + assert tracker._total_energy.kWh > 0 + finally: + tracker.stop() + assert samples == [] + + +def test_conservation_with_fake_clocks_and_energy(): + """Random windows, meters and GPU load: nothing is created or lost.""" + import random + + rng = random.Random(1428) + times = _FakeTimes() + attributor = EnergyAttributor(cpu_times=times) + gpu = dict(gpu_quality="measured", cpu_idle_w=None) + totals = dict(cpu=0.0, gpu=0.0, ram=0.0) + attributor.reset_window(_sample(0.0, timestamp=0.0, **gpu)) + results = [] + live = [] + for t in range(1, 500): + for _ in range(rng.randint(0, 4)): + state = _begin(attributor, start=t - rng.random()) + state.on_resolved = results.append + live.append(state) + for state in live: + state.meter.ns += rng.randint(0, 50_000_000) + for state in rng.sample(live, k=len(live) // 2): + attributor.end(state) + live.remove(state) + for part, watts in (("cpu", 80), ("gpu", 250), ("ram", 5)): + totals[part] += _kwh(watts * rng.random(), 1) + times.advance(process=rng.random(), busy=rng.random() * 4, total=8.0) + attributor.on_window( + _sample( + sum(totals.values()), + cpu=totals["cpu"], + gpu=totals["gpu"], + ram_kwh=totals["ram"], + timestamp=float(t), + **gpu, + ) + ) + _invariant(attributor) + attributor.close() + report = attributor.report() + assert report["settled_kwh"] == pytest.approx(sum(totals.values()), rel=1e-12) + charged = sum(r.energy_kwh + r.gpu_kwh for r in results if r.windows) + assert charged == pytest.approx(report["attributed_kwh"], rel=1e-12) + assert all(report[k] > 0 for k in ("idle_kwh", "other_processes_kwh")) + + +def _metered_app(seen): + app = FastAPI() + app.add_middleware( + CodeCarbonMiddleware, + tracker=_FakeTracker(), + on_request=lambda energy, kg, status: seen.append(energy), + ) + + @app.get("/burn/{ms}") + async def burn(ms: float): + _burn(ms / 1000) + return {} + + @app.get("/burn-sync/{ms}") + def burn_sync(ms: float): + _burn(ms / 1000) + return {} + + @app.get("/sleep/{ms}") + async def sleep(ms: float): + await asyncio.sleep(ms / 1000) + return {} + + return app + + +def _cpu_seconds(path: str, repeat: int = 5) -> list[float]: + seen = [] + with TestClient(_metered_app(seen)) as client: + client.get(path) # warm-up: first-call imports and caches + for _ in range(repeat): + assert client.get(path).status_code == 200 + # Lifespan shutdown closed the middleware and emitted everything. + return [energy.cpu_seconds for energy in seen[1:]] + + +def test_cpu_seconds_of_a_known_async_burn(): + for cpu_s in _cpu_seconds("/burn/20"): + assert cpu_s == pytest.approx(0.020, rel=0.10, abs=0.0002) + + +def test_cpu_seconds_of_a_known_sync_burn(): + for cpu_s in _cpu_seconds("/burn-sync/20"): + assert cpu_s == pytest.approx(0.020, rel=0.10, abs=0.0002) + + +def _threadpool_app(seen, **middleware): + """Every way FastAPI and Starlette send a request's work to a thread.""" + from fastapi import APIRouter, Depends + from fastapi.responses import StreamingResponse + + app = FastAPI() + app.add_middleware( + CodeCarbonMiddleware, + tracker=_FakeTracker(), + on_request=lambda energy, kg, status: seen.append(energy), + **middleware, + ) + + def burn_dep() -> int: + _burn(0.010) + return 1 + + def real_value() -> str: + return "real" + + @app.get("/sync-dep") + async def sync_dep(x: int = Depends(burn_dep)): + return {} + + @app.get("/stream") + def stream(): + def chunks(): + for _ in range(4): + _burn(0.005) + yield b"x" + + return StreamingResponse(chunks()) + + router = APIRouter() + + @router.get("/sync") + def included_sync(): + _burn(0.010) + return {} + + app.include_router(router, prefix="/included") + + @app.get("/value") + def value(v: str = Depends(real_value)): + return {"v": v} + + app.state.real_value = real_value + return app + + +@pytest.mark.parametrize( + "path, cpu_s", + [("/sync-dep", 0.010), ("/stream", 0.020), ("/included/sync", 0.010)], +) +def test_threadpool_work_is_metered(path, cpu_s): + seen = [] + with TestClient(_threadpool_app(seen)) as client: + client.get(path) # warm-up + for _ in range(3): + assert client.get(path).status_code == 200 + for energy in seen[1:]: + assert energy.cpu_seconds == pytest.approx(cpu_s, rel=0.10, abs=0.0005) + + +def test_plain_starlette_sync_endpoint_is_metered(): + from starlette.applications import Starlette + from starlette.middleware import Middleware + from starlette.responses import PlainTextResponse + from starlette.routing import Route + + def endpoint(request): + _burn(0.010) + return PlainTextResponse("ok") + + seen = [] + app = Starlette( + routes=[Route("/", endpoint)], + middleware=[ + Middleware( + CodeCarbonMiddleware, + tracker=_FakeTracker(), + on_request=lambda energy, kg, status: seen.append(energy), + ) + ], + ) + with TestClient(app) as client: + for _ in range(3): + assert client.get("/").status_code == 200 + for energy in seen: + assert energy.cpu_seconds == pytest.approx(0.010, rel=0.10, abs=0.0005) + + +def test_dependency_overrides_still_work(): + seen = [] + app = _threadpool_app(seen) + app.dependency_overrides[app.state.real_value] = lambda: "override" + with TestClient(app) as client: + assert client.get("/value").json() == {"v": "override"} + + +def test_threadpool_metering_can_be_turned_off(): + import anyio.to_thread + + original = anyio.to_thread.run_sync + seen = [] + with TestClient(_threadpool_app(seen, meter_threadpool=False)) as client: + client.get("/included/sync") + assert anyio.to_thread.run_sync is original + assert seen[0].cpu_seconds < 0.005 + + +def test_run_sync_patch_is_transparent_and_restored(): + import anyio + import anyio.to_thread + + from codecarbon.integrations.fastapi import middleware as mw + + original = anyio.to_thread.run_sync + limiter = anyio.CapacityLimiter(1) + + def add(a, b): + return a + b + + def boom(): + raise ValueError("boom") + + async def main(): + # No request meter set: plain pass-through, arguments and errors. + assert await anyio.to_thread.run_sync(add, 1, 2, limiter=limiter) == 3 + with pytest.raises(ValueError): + await anyio.to_thread.run_sync(boom) + meter = _Meter() + token = mw._current_meter.set(meter) + try: + await anyio.to_thread.run_sync(_burn, 0.005) + finally: + mw._current_meter.reset(token) + return meter + + mw._patch_run_sync() + mw._patch_run_sync() # idempotent: one wrapper, two users + try: + assert anyio.to_thread.run_sync is mw._metered_run_sync + meter = anyio.run(main) + assert meter.total_ns() / 1e9 == pytest.approx(0.005, rel=0.2) + mw._unpatch_run_sync() + assert anyio.to_thread.run_sync is mw._metered_run_sync + finally: + mw._unpatch_run_sync() + assert anyio.to_thread.run_sync is original + + +def test_threadpool_patch_restored_on_shutdown(): + import anyio.to_thread + + original = anyio.to_thread.run_sync + with TestClient(_threadpool_app([])) as client: + client.get("/included/sync") + assert anyio.to_thread.run_sync is not original + assert anyio.to_thread.run_sync is original + + +@pytest.mark.skipif( + not hasattr(time, "pthread_getcpuclockid"), + reason="no way to read another thread's CPU clock on this platform", +) +def test_long_sync_call_is_charged_while_running(): + meter = _Meter() + started = threading.Event() + + def work(): + started.set() + _burn(0.2) + + worker = threading.Thread(target=meter.run_in_thread, args=(work,)) + worker.start() + started.wait() + time.sleep(0.1) + midway = meter.total_ns() + worker.join() + assert 0.02 < midway / 1e9 < 0.2 + assert meter.total_ns() / 1e9 == pytest.approx(0.2, rel=0.1) + + +def test_sleeping_endpoint_uses_almost_no_cpu(): + for cpu_s in _cpu_seconds("/sleep/50"): + assert cpu_s < 0.002