From 7b566a8f1fc46c29f715d784f6eb5a710a7b069c Mon Sep 17 00:00:00 2001 From: David Berenstein Date: Thu, 20 Aug 2026 08:12:57 +0200 Subject: [PATCH 1/4] feat(api): reuse connections, add timeouts, and stamp emissions with measurement time `ApiClient` called the module-level `requests` functions, so every call opened a new TCP connection and TLS handshake. A tracker sending one measurement per tick opened 50 connections for 50 uploads; with a `Session` it opens 1. The flat 2s timeout is replaced with `(3.05, 10)`: the old value timed out against a healthy but loaded API, while a hung endpoint still cannot block the scheduler thread for long. `ApiClient.add_emission` also discarded `carbon_emission["timestamp"]` and called `get_datetime_with_timezone()` instead, so every row stored the moment the payload was built rather than the moment it was measured. Harmless while the two are milliseconds apart, wrong by the full latency as soon as a send is slow, retried or queued. The measurement timestamp is now normalised to offset-aware, falling back to now when the payload carries none or an unparseable one (the method is public and takes a plain dict). CSV output is untouched. Co-Authored-By: Claude Opus 5 (1M context) --- codecarbon/core/api_client.py | 62 +++++++++++---- codecarbon/output_methods/http.py | 3 + tests/test_api_call.py | 40 ++++++++++ tests/test_api_client_session.py | 126 ++++++++++++++++++++++++++++++ 4 files changed, 214 insertions(+), 17 deletions(-) create mode 100644 tests/test_api_client_session.py diff --git a/codecarbon/core/api_client.py b/codecarbon/core/api_client.py index bc2e0974e..450cdabbc 100644 --- a/codecarbon/core/api_client.py +++ b/codecarbon/core/api_client.py @@ -8,7 +8,7 @@ # from httpx import AsyncClient import dataclasses import json -from datetime import timedelta, tzinfo +from datetime import datetime, timedelta, tzinfo import requests @@ -33,6 +33,26 @@ def get_datetime_with_timezone(): return str(arrow.now().isoformat()) +# (connect, read) seconds, replacing a flat 2s that timed out on a loaded API. +_TIMEOUT = (3.05, 10) + + +def _measurement_timestamp(carbon_emission: dict) -> str: + """ + Offset-aware ISO timestamp of *when the measurement was taken*, taken from + EmissionsData.timestamp. Falls back to now for hand-built payloads that + carry no usable timestamp. + """ + try: + return ( + datetime.fromisoformat(carbon_emission["timestamp"]) + .astimezone() + .isoformat() + ) + except (KeyError, TypeError, ValueError): + return get_datetime_with_timezone() + + class ApiClient: # (AsyncClient) """ This class call the Code Carbon API @@ -58,6 +78,8 @@ def __init__( :create_run_automatically: If False, do not create a run. To use API in read only mode. """ # super().__init__(base_url=endpoint_url) # (AsyncClient) + # A Session so the socket and TLS handshake are reused across calls. + self._session = requests.Session() self.url = endpoint_url self.experiment_id = experiment_id self.api_key = api_key @@ -80,16 +102,20 @@ def _request(self, method, url, payload=None, expected_status=200): Call the API and return the response, raising on anything that is not the status code the API answers on success. - :method: the requests function to call, for example requests.get + :method: the session function to call, for example self._session.get :payload: the JSON body to send, if any :expected_status: the http code the API returns when the call succeeds """ headers = self._get_headers() - response = method(url=url, json=payload, timeout=2, headers=headers) + response = method(url=url, json=payload, timeout=_TIMEOUT, headers=headers) if response.status_code != expected_status: self._raise_api_error(url, payload or {}, response) return response + def close(self): + """Release the pooled sockets. Safe to call more than once.""" + self._session.close() + def set_access_token(self, token: str): """This method sets the access token to be used for the API. Args: @@ -102,14 +128,14 @@ def check_auth(self): Check API access to user account """ url = self.url + "/auth/check" - return self._request(requests.get, url).json() + return self._request(self._session.get, url).json() def get_list_organizations(self): """ List all organizations """ url = self.url + "/organizations" - return self._request(requests.get, url).json() + return self._request(self._session.get, url).json() def check_organization_exists(self, organization_name: str): """ @@ -134,7 +160,7 @@ def create_organization(self, organization: OrganizationCreate): return organization else: return self._request( - requests.post, url, payload=payload, expected_status=201 + self._session.post, url, payload=payload, expected_status=201 ).json() def get_organization(self, organization_id): @@ -142,7 +168,7 @@ def get_organization(self, organization_id): Get an organization """ url = self.url + "/organizations/" + organization_id - return self._request(requests.get, url).json() + return self._request(self._session.get, url).json() def update_organization(self, organization: OrganizationCreate): """ @@ -150,14 +176,14 @@ def update_organization(self, organization: OrganizationCreate): """ payload = dataclasses.asdict(organization) url = self.url + "/organizations/" + organization.id - return self._request(requests.patch, url, payload=payload).json() + return self._request(self._session.patch, url, payload=payload).json() def list_projects_from_organization(self, organization_id): """ List all projects """ url = self.url + "/organizations/" + organization_id + "/projects" - return self._request(requests.get, url).json() + return self._request(self._session.get, url).json() def create_project(self, project: ProjectCreate): """ @@ -166,7 +192,7 @@ def create_project(self, project: ProjectCreate): payload = dataclasses.asdict(project) url = self.url + "/projects" return self._request( - requests.post, url, payload=payload, expected_status=201 + self._session.post, url, payload=payload, expected_status=201 ).json() def get_project(self, project_id): @@ -174,7 +200,7 @@ def get_project(self, project_id): Get a project """ url = self.url + "/projects/" + project_id - return self._request(requests.get, url).json() + return self._request(self._session.get, url).json() def add_emission(self, carbon_emission: dict): assert self.experiment_id is not None @@ -195,7 +221,7 @@ def add_emission(self, carbon_emission: dict): ) return False emission = EmissionCreate( - timestamp=get_datetime_with_timezone(), + timestamp=_measurement_timestamp(carbon_emission), run_id=self.run_id, duration=int(carbon_emission["duration"]), emissions_sum=carbon_emission["emissions"], @@ -215,7 +241,7 @@ def add_emission(self, carbon_emission: dict): try: payload = dataclasses.asdict(emission) url = self.url + "/emissions" - self._request(requests.post, url, payload=payload, expected_status=201) + self._request(self._session.post, url, payload=payload, expected_status=201) logger.debug(f"ApiClient - Successful upload emission {payload} to {url}") except requests.exceptions.HTTPError: # Already logged by _raise_api_error, do not log it twice. @@ -256,7 +282,9 @@ def _create_run(self, experiment_id: str): ) payload = dataclasses.asdict(run) url = self.url + "/runs" - r = self._request(requests.post, url, payload=payload, expected_status=201) + r = self._request( + self._session.post, url, payload=payload, expected_status=201 + ) self.run_id = r.json()["id"] logger.info( "ApiClient Successfully registered your run on the API.\n\n" @@ -282,7 +310,7 @@ def list_experiments_from_project(self, project_id: str): List all experiments for a project """ url = self.url + "/projects/" + project_id + "/experiments" - return self._request(requests.get, url).json() + return self._request(self._session.get, url).json() def set_experiment(self, experiment_id: str): """ @@ -298,7 +326,7 @@ def add_experiment(self, experiment: ExperimentCreate): payload = dataclasses.asdict(experiment) url = self.url + "/experiments" return self._request( - requests.post, url, payload=payload, expected_status=201 + self._session.post, url, payload=payload, expected_status=201 ).json() def get_experiment(self, experiment_id): @@ -306,7 +334,7 @@ def get_experiment(self, experiment_id): Get an experiment by id """ url = self.url + "/experiments/" + experiment_id - return self._request(requests.get, url).json() + return self._request(self._session.get, url).json() def _raise_api_error(self, url, payload, response): """ diff --git a/codecarbon/output_methods/http.py b/codecarbon/output_methods/http.py index e0ff710b1..b7e8a5896 100644 --- a/codecarbon/output_methods/http.py +++ b/codecarbon/output_methods/http.py @@ -57,6 +57,9 @@ def __init__( ) self.run_id = self.api.run_id + def exit(self) -> None: + self.api.close() + def _ensure_api_run(self) -> None: if self.api.run_id is None and self.api.experiment_id is not None: self.api._create_run(self.api.experiment_id) diff --git a/tests/test_api_call.py b/tests/test_api_call.py index 31e25c039..481c76111 100644 --- a/tests/test_api_call.py +++ b/tests/test_api_call.py @@ -1,5 +1,6 @@ import dataclasses import unittest +from datetime import datetime from uuid import uuid4 import requests @@ -261,6 +262,45 @@ def test_add_emission_skips_short_duration(self): ) ) + def test_add_emission_keeps_measurement_timestamp(self): + """The row must carry when it was measured, not when it was sent.""" + payload = { + "duration": 10, + "emissions": 1.0, + "emissions_rate": 1.0, + "cpu_power": 1.0, + "gpu_power": 0.0, + "ram_power": 0.5, + "cpu_energy": 0.1, + "gpu_energy": 0.0, + "ram_energy": 0.1, + "energy_consumed": 0.2, + } + with requests_mock.Mocker() as m: + m.post("http://test.com/emissions", status_code=201) + api = ApiClient( + endpoint_url="http://test.com", + experiment_id="exp-1", + conf=conf, + create_run_automatically=False, + ) + api.run_id = "run-1" + + # naive timestamp, as produced by EmissionsData + assert api.add_emission({**payload, "timestamp": "2020-01-01T00:00:00"}) + sent = datetime.fromisoformat(m.last_request.json()["timestamp"]) + self.assertEqual( + sent.replace(tzinfo=None).isoformat(), "2020-01-01T00:00:00" + ) + self.assertIsNotNone(sent.tzinfo) + + # missing / unparseable timestamps fall back to now + for bad in ({}, {"timestamp": None}, {"timestamp": "222"}): + assert api.add_emission({**payload, **bad}) + sent = datetime.fromisoformat(m.last_request.json()["timestamp"]) + self.assertIsNotNone(sent.tzinfo) + self.assertGreater(sent.year, 2020) + def test_add_emission_raises_on_unsuccessful_post(self): with requests_mock.Mocker() as m: m.post("http://test.com/emissions", text="bad", status_code=500) diff --git a/tests/test_api_client_session.py b/tests/test_api_client_session.py new file mode 100644 index 000000000..085881785 --- /dev/null +++ b/tests/test_api_client_session.py @@ -0,0 +1,126 @@ +""" +Connection-reuse and timeout tests for ApiClient. + +These run against a stdlib HTTP server on loopback rather than requests_mock, +because requests_mock replaces the transport adapter and therefore never opens +a real connection, which is exactly what is under test here. No traffic leaves +the machine. +""" + +import threading +import unittest +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer + +from codecarbon.core.api_client import ApiClient + +CONF = { + "os": "linux", + "python_version": "3.12", + "codecarbon_version": "3.0", + "cpu_count": 8, + "cpu_model": "CPU", + "gpu_count": 0, + "gpu_model": "", + "longitude": 0.0, + "latitude": 0.0, + "region": "EU", + "provider": "none", + "ram_total_size": 16.0, + "tracking_mode": "machine", +} + +EMISSION = { + "duration": 5, + "emissions": 1.0, + "emissions_rate": 1.0, + "cpu_power": 1.0, + "gpu_power": 0.0, + "ram_power": 0.5, + "cpu_energy": 0.1, + "gpu_energy": 0.0, + "ram_energy": 0.1, + "energy_consumed": 0.2, +} + + +class _Handler(BaseHTTPRequestHandler): + protocol_version = "HTTP/1.1" # keep-alive, so pooling is observable + + def log_message(self, *args): + pass + + def _serve(self): + self.server.state["requests"] += 1 + length = int(self.headers.get("Content-Length", 0) or 0) + if length: + self.rfile.read(length) + body = b'{"id": "run-1"}' + self.send_response(201) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(body))) + self.end_headers() + self.wfile.write(body) + + do_GET = _serve + do_POST = _serve + + +class _Server(ThreadingHTTPServer): + daemon_threads = True + allow_reuse_address = True + + def __init__(self, state): + self.state = state + super().__init__(("127.0.0.1", 0), _Handler) + + def process_request(self, request, client_address): + self.state["connections"] += 1 + super().process_request(request, client_address) + + +class TestSessionReuse(unittest.TestCase): + def setUp(self): + self.state = {"requests": 0, "connections": 0} + self.server = _Server(self.state) + threading.Thread(target=self.server.serve_forever, daemon=True).start() + self.url = f"http://127.0.0.1:{self.server.server_address[1]}" + self.addCleanup(self.server.server_close) + self.addCleanup(self.server.shutdown) + self.api = ApiClient( + endpoint_url=self.url, + experiment_id="exp-1", + conf=CONF, + create_run_automatically=False, + ) + self.addCleanup(self.api.close) + self.api.run_id = "run-1" + + def test_sequential_calls_reuse_one_connection(self): + for _ in range(50): + self.assertTrue(self.api.add_emission(dict(EMISSION))) + + self.assertEqual(self.state["requests"], 50) + self.assertEqual(self.state["connections"], 1) + + def test_close_is_idempotent(self): + self.api.add_emission(dict(EMISSION)) + self.api.close() + self.api.close() + + +class TestTimeout(unittest.TestCase): + def test_requests_get_a_connect_and_read_timeout(self): + api = ApiClient(endpoint_url="http://test.com", create_run_automatically=False) + self.addCleanup(api.close) + seen = {} + + def fake_get(url, json, timeout, headers): + seen["timeout"] = timeout + return type("R", (), {"status_code": 200, "json": lambda self: {}})() + + api._request(fake_get, "http://test.com/x") + self.assertEqual(seen["timeout"], (3.05, 10)) + + +if __name__ == "__main__": + unittest.main() From 72c53b779f538695400364b33eaf0d02c9202462 Mon Sep 17 00:00:00 2001 From: David Berenstein Date: Wed, 23 Sep 2026 16:42:58 +0900 Subject: [PATCH 2/4] fix(api): cool down after a failed run creation A down API now costs one blocking run-creation call per minute instead of one per measurement tick. Co-Authored-By: Claude Opus 5.5 (1M context) --- codecarbon/core/api_client.py | 14 ++++++++++++++ tests/test_api_call.py | 20 ++++++++++++++++++++ 2 files changed, 34 insertions(+) diff --git a/codecarbon/core/api_client.py b/codecarbon/core/api_client.py index 450cdabbc..f6dd0f93f 100644 --- a/codecarbon/core/api_client.py +++ b/codecarbon/core/api_client.py @@ -8,6 +8,7 @@ # from httpx import AsyncClient import dataclasses import json +import time from datetime import datetime, timedelta, tzinfo import requests @@ -35,6 +36,9 @@ def get_datetime_with_timezone(): # (connect, read) seconds, replacing a flat 2s that timed out on a loaded API. _TIMEOUT = (3.05, 10) +# Seconds to wait after a failed run creation before trying again, so a down +# API costs one blocking call per minute instead of one per measurement. +_RUN_CREATE_COOLDOWN = 60 def _measurement_timestamp(carbon_emission: dict) -> str: @@ -85,6 +89,7 @@ def __init__( self.api_key = api_key self.conf = conf self.access_token = access_token + self._run_create_failed_at = None if self.experiment_id is not None and create_run_automatically: self._create_run(self.experiment_id) @@ -261,6 +266,14 @@ def _create_run(self, experiment_id: str): "ApiClient FATAL The ApiClient._create_run() needs an experiment_id !" ) return None + if ( + self._run_create_failed_at is not None + and time.monotonic() - self._run_create_failed_at < _RUN_CREATE_COOLDOWN + ): + logger.debug("ApiClient - run creation failed recently, not retrying yet") + return None + # Cleared on success below; set now so every failure path is covered. + self._run_create_failed_at = time.monotonic() try: run = RunCreate( timestamp=get_datetime_with_timezone(), @@ -286,6 +299,7 @@ def _create_run(self, experiment_id: str): self._session.post, url, payload=payload, expected_status=201 ) self.run_id = r.json()["id"] + self._run_create_failed_at = None logger.info( "ApiClient Successfully registered your run on the API.\n\n" + f"Run ID: {self.run_id}\n" diff --git a/tests/test_api_call.py b/tests/test_api_call.py index 481c76111..e4e0b8835 100644 --- a/tests/test_api_call.py +++ b/tests/test_api_call.py @@ -343,6 +343,26 @@ def test_create_run_raises_on_unsuccessful_status(self): api._create_run("experiment_id") self.assertIsNone(api.run_id) + def test_failed_create_run_is_not_retried_during_cooldown(self): + with requests_mock.Mocker() as m: + runs = m.post("http://test.com/runs", text="down", status_code=503) + api = ApiClient( + endpoint_url="http://test.com", + experiment_id="experiment_id", + api_key="Toto", + conf=conf, + create_run_automatically=False, + ) + with self.assertRaises(requests.exceptions.HTTPError): + api._create_run("experiment_id") + self.assertIsNone(api._create_run("experiment_id")) + self.assertEqual(runs.call_count, 1) + + api._run_create_failed_at -= 3600 # cooldown elapsed + runs = m.post("http://test.com/runs", json={"id": "run-1"}, status_code=201) + self.assertEqual(api._create_run("experiment_id"), "run-1") + self.assertIsNone(api._run_create_failed_at) + def test_create_run_raises_on_unexpected_2xx_status(self): with requests_mock.Mocker() as m: m.post("http://test.com/runs", json={}, status_code=200) From a8fce86e066a6cf19460eca2057d478da1aa0eec Mon Sep 17 00:00:00 2001 From: David Berenstein Date: Mon, 28 Sep 2026 09:02:51 +0900 Subject: [PATCH 3/4] fix(api): let the final flush bypass the run-creation cooldown --- codecarbon/core/api_client.py | 37 +++++++--- codecarbon/emissions_tracker.py | 8 +++ codecarbon/output_methods/http.py | 14 ++-- tests/output_methods/test_http.py | 25 ++++++- tests/test_api_call.py | 109 +++++++++++++++++++++++++++++- 5 files changed, 175 insertions(+), 18 deletions(-) diff --git a/codecarbon/core/api_client.py b/codecarbon/core/api_client.py index f6dd0f93f..9b89d7527 100644 --- a/codecarbon/core/api_client.py +++ b/codecarbon/core/api_client.py @@ -207,19 +207,35 @@ def get_project(self, project_id): url = self.url + "/projects/" + project_id return self._request(self._session.get, url).json() - def add_emission(self, carbon_emission: dict): + def add_emission(self, carbon_emission: dict, final: bool = False): assert self.experiment_id is not None if self.run_id is None: + # Captured before the call: tells us whether _create_run is about + # to skip its attempt because of the cooldown, so we can log the + # expected per-tick skip at debug instead of as an error. + in_cooldown = ( + not final + and self._run_create_failed_at is not None + and time.monotonic() - self._run_create_failed_at < _RUN_CREATE_COOLDOWN + ) logger.warning( "ApiClient.add_emission() need a run_id : the initial call may " + "have failed. Retrying..." ) - self._create_run(self.experiment_id) + self._create_run(self.experiment_id, bypass_cooldown=final) if self.run_id is None: - logger.error( - "ApiClient.add_emission still no run_id, aborting for this time !" - ) - return False + if in_cooldown: + logger.debug( + "ApiClient.add_emission still no run_id, run creation" + " is in its cooldown, will retry later." + ) + else: + logger.error( + "ApiClient.add_emission still no run_id, aborting for this time !" + ) + return False + # Run creation just succeeded (e.g. a final-flush bypass): fall + # through and send this emission instead of dropping it. if carbon_emission["duration"] < 1: logger.warning( "ApiClient : emissions not sent because of a duration smaller than 1." @@ -256,9 +272,13 @@ def add_emission(self, carbon_emission: dict): raise return True - def _create_run(self, experiment_id: str): + def _create_run(self, experiment_id: str, bypass_cooldown: bool = False): """ Create the experiment for project_id + + :bypass_cooldown: skip the cooldown check and retry immediately, used + by the final flush on tracker stop/exit so the last emission is + not silently dropped because of a recent failure. """ if self.experiment_id is None: # TODO : raise an Exception ? @@ -267,7 +287,8 @@ def _create_run(self, experiment_id: str): ) return None if ( - self._run_create_failed_at is not None + not bypass_cooldown + and self._run_create_failed_at is not None and time.monotonic() - self._run_create_failed_at < _RUN_CREATE_COOLDOWN ): logger.debug("ApiClient - run creation failed recently, not retrying yet") diff --git a/codecarbon/emissions_tracker.py b/codecarbon/emissions_tracker.py index 96ed00c91..f384028c0 100644 --- a/codecarbon/emissions_tracker.py +++ b/codecarbon/emissions_tracker.py @@ -32,6 +32,7 @@ from codecarbon.lock import Lock from codecarbon.output_methods.base_output import BaseOutput, OutputMethod from codecarbon.output_methods.emissions_data import EmissionsData +from codecarbon.output_methods.http import CodeCarbonAPIOutput if TYPE_CHECKING: from codecarbon.external.geography import CloudMetadata, GeoMetadata @@ -931,6 +932,13 @@ def stop(self) -> Optional[float]: experiment_name=self._experiment_name, ) + # If run creation was still in its cooldown, the emission above was + # dropped. This is the last chance to send it, so bypass the + # cooldown and retry once instead of losing the row. + for handler in self._output_handlers: + if isinstance(handler, CodeCarbonAPIOutput) and handler.run_id is None: + handler.out(emissions_data, emissions_data_delta, final=True) + self.final_emissions_data = emissions_data self.final_emissions = emissions_data.emissions diff --git a/codecarbon/output_methods/http.py b/codecarbon/output_methods/http.py index b7e8a5896..bc2866641 100644 --- a/codecarbon/output_methods/http.py +++ b/codecarbon/output_methods/http.py @@ -60,20 +60,20 @@ def __init__( def exit(self) -> None: self.api.close() - def _ensure_api_run(self) -> None: + def _ensure_api_run(self, final: bool = False) -> None: if self.api.run_id is None and self.api.experiment_id is not None: - self.api._create_run(self.api.experiment_id) + self.api._create_run(self.api.experiment_id, bypass_cooldown=final) self.run_id = self.api.run_id - def _emit(self, delta: EmissionsData) -> None: + def _emit(self, delta: EmissionsData, final: bool = False) -> None: try: - self._ensure_api_run() - self.api.add_emission(dataclasses.asdict(delta)) + self._ensure_api_run(final=final) + self.api.add_emission(dataclasses.asdict(delta), final=final) except Exception as e: logger.error(e, exc_info=True) def live_out(self, _, delta: EmissionsData): self._emit(delta) - def out(self, _, delta: EmissionsData): - self._emit(delta) + def out(self, _, delta: EmissionsData, final: bool = False): + self._emit(delta, final=final) diff --git a/tests/output_methods/test_http.py b/tests/output_methods/test_http.py index 56d909b46..5c2541394 100644 --- a/tests/output_methods/test_http.py +++ b/tests/output_methods/test_http.py @@ -175,14 +175,14 @@ def test_codecarbon_api_live_out_creates_run_when_missing(self): ) api_output.api.run_id = None - def create_run(experiment_id): + def create_run(experiment_id, bypass_cooldown=False): api_output.api.run_id = "run-created" return "run-created" mock_create_run.side_effect = create_run api_output.live_out(None, self.emissions_data) - mock_create_run.assert_called_once_with("exp-1") + mock_create_run.assert_called_once_with("exp-1", bypass_cooldown=False) self.assertEqual(api_output.api.run_id, "run-created") self.assertEqual(api_output.run_id, "run-created") @@ -209,6 +209,27 @@ def test_codecarbon_api_out(self): api_output.out(None, self.emissions_data) self.mock_add_emission.assert_called_once() + def test_codecarbon_api_out_final_bypasses_cooldown(self): + """The final flush (tracker.stop()) must bypass the run-creation + cooldown and mark the emission as final, so it is not dropped when a + run was recently failing to be created.""" + with patch( + "codecarbon.output_methods.http.ApiClient._create_run" + ) as mock_create_run: + api_output = CodeCarbonAPIOutput( + endpoint_url=self.url, + experiment_id="exp-1", + api_key=self.api_key, + conf=None, + ) + api_output.api.run_id = None + + api_output.out(None, self.emissions_data, final=True) + + mock_create_run.assert_called_once_with("exp-1", bypass_cooldown=True) + self.mock_add_emission.assert_called_once() + self.assertTrue(self.mock_add_emission.call_args.kwargs.get("final")) + @patch("codecarbon.output_methods.http.logger.error") def test_codecarbon_out_api_call_failure(self, mock_logger): self.mock_add_emission.side_effect = Exception("Test exception") diff --git a/tests/test_api_call.py b/tests/test_api_call.py index e4e0b8835..ac042128f 100644 --- a/tests/test_api_call.py +++ b/tests/test_api_call.py @@ -217,7 +217,7 @@ def test_add_emission_returns_false_when_run_creation_fails(self): create_run_automatically=False, ) - api._create_run = lambda experiment_id: None + api._create_run = lambda experiment_id, bypass_cooldown=False: None self.assertFalse( api.add_emission( @@ -363,6 +363,113 @@ def test_failed_create_run_is_not_retried_during_cooldown(self): self.assertEqual(api._create_run("experiment_id"), "run-1") self.assertIsNone(api._run_create_failed_at) + def test_create_run_bypasses_cooldown_when_requested(self): + with requests_mock.Mocker() as m: + m.post("http://test.com/runs", text="down", status_code=503) + api = ApiClient( + endpoint_url="http://test.com", + experiment_id="experiment_id", + api_key="Toto", + conf=conf, + create_run_automatically=False, + ) + with self.assertRaises(requests.exceptions.HTTPError): + api._create_run("experiment_id") + + m.post("http://test.com/runs", json={"id": "run-1"}, status_code=201) + # Still well within the cooldown, but the caller asked to bypass it. + self.assertEqual( + api._create_run("experiment_id", bypass_cooldown=True), "run-1" + ) + + def test_add_emission_logs_cooldown_skip_at_debug_not_error(self): + with requests_mock.Mocker() as m: + m.post("http://test.com/runs", text="down", status_code=503) + api = ApiClient( + endpoint_url="http://test.com", + experiment_id="exp-1", + api_key="Toto", + conf=conf, + create_run_automatically=False, + ) + # First call: the run-creation attempt itself is a genuine + # failure (raises, same as today). + with self.assertRaises(requests.exceptions.HTTPError): + api.add_emission( + { + "duration": 2, + "emissions": 1.0, + "emissions_rate": 1.0, + "cpu_power": 1.0, + "gpu_power": 0.0, + "ram_power": 0.5, + "cpu_energy": 0.1, + "gpu_energy": 0.0, + "ram_energy": 0.1, + "energy_consumed": 0.2, + } + ) + + with self.assertLogs("codecarbon", level="DEBUG") as cm2: + self.assertFalse( + api.add_emission( + { + "duration": 2, + "emissions": 1.0, + "emissions_rate": 1.0, + "cpu_power": 1.0, + "gpu_power": 0.0, + "ram_power": 0.5, + "cpu_energy": 0.1, + "gpu_energy": 0.0, + "ram_energy": 0.1, + "energy_consumed": 0.2, + } + ) + ) + # Second call: still within the cooldown, so no new attempt is + # made and the per-tick skip is only a DEBUG log, not ERROR. + second_errors = [r for r in cm2.records if r.levelname == "ERROR"] + self.assertFalse(second_errors) + + def test_add_emission_final_bypasses_cooldown_and_sends(self): + with requests_mock.Mocker() as m: + m.post("http://test.com/runs", text="down", status_code=503) + api = ApiClient( + endpoint_url="http://test.com", + experiment_id="exp-1", + api_key="Toto", + conf=conf, + create_run_automatically=False, + ) + emission = { + "duration": 2, + "emissions": 1.0, + "emissions_rate": 1.0, + "cpu_power": 1.0, + "gpu_power": 0.0, + "ram_power": 0.5, + "cpu_energy": 0.1, + "gpu_energy": 0.0, + "ram_energy": 0.1, + "energy_consumed": 0.2, + } + # First tick: run creation fails (raises) and enters cooldown. + with self.assertRaises(requests.exceptions.HTTPError): + api.add_emission(emission) + + # The server recovers, but we are still inside the cooldown window. + m.post("http://test.com/runs", json={"id": "run-1"}, status_code=201) + m.post("http://test.com/emissions", json={"id": "em-1"}, status_code=201) + + # A normal (non-final) tick still respects the cooldown. + self.assertFalse(api.add_emission(emission)) + self.assertIsNone(api.run_id) + + # The final flush bypasses the cooldown and retries once. + self.assertTrue(api.add_emission(emission, final=True)) + self.assertEqual(api.run_id, "run-1") + def test_create_run_raises_on_unexpected_2xx_status(self): with requests_mock.Mocker() as m: m.post("http://test.com/runs", json={}, status_code=200) From 38be0d978bb001f039bda7c50c628c434f6378c9 Mon Sep 17 00:00:00 2001 From: David Berenstein Date: Mon, 28 Sep 2026 12:29:59 +0900 Subject: [PATCH 4/4] fix(api): skip the final run-creation retry when stop() already tried --- codecarbon/emissions_tracker.py | 15 +++++++++--- tests/test_emissions_tracker.py | 41 +++++++++++++++++++++++++++++++++ 2 files changed, 53 insertions(+), 3 deletions(-) diff --git a/codecarbon/emissions_tracker.py b/codecarbon/emissions_tracker.py index f384028c0..de638dacb 100644 --- a/codecarbon/emissions_tracker.py +++ b/codecarbon/emissions_tracker.py @@ -926,6 +926,7 @@ def stop(self) -> Optional[float]: emissions_data = self._prepare_emissions_data() emissions_data_delta = self._compute_emissions_delta(emissions_data) + persist_started_at = time.monotonic() self._persist_data( total_emissions=emissions_data, delta_emissions=emissions_data_delta, @@ -933,11 +934,19 @@ def stop(self) -> Optional[float]: ) # If run creation was still in its cooldown, the emission above was - # dropped. This is the last chance to send it, so bypass the - # cooldown and retry once instead of losing the row. + # dropped without even trying the API. This is the last chance to + # send it, so bypass the cooldown and retry once instead of losing + # the row. But if the persist call above already made (and failed) a + # run-creation attempt, retrying again here would just hit the + # server a second time for nothing. for handler in self._output_handlers: if isinstance(handler, CodeCarbonAPIOutput) and handler.run_id is None: - handler.out(emissions_data, emissions_data_delta, final=True) + already_attempted = ( + handler.api._run_create_failed_at is not None + and handler.api._run_create_failed_at >= persist_started_at + ) + if not already_attempted: + handler.out(emissions_data, emissions_data_delta, final=True) self.final_emissions_data = emissions_data self.final_emissions = emissions_data.emissions diff --git a/tests/test_emissions_tracker.py b/tests/test_emissions_tracker.py index 8ab12e5d8..68e8f5cd9 100644 --- a/tests/test_emissions_tracker.py +++ b/tests/test_emissions_tracker.py @@ -1108,3 +1108,44 @@ def test_cumulative_emissions_with_varying_intensity( # Verification: If it wasn't cumulative, it would be 3.0 kWh * 300 g/kWh = 0.9 kg self.assertLess(data3.emissions, 0.8) + + @responses.activate + def test_stop_does_not_retry_run_creation_twice_when_server_down( + self, + mock_cli_setup, + mock_log_values, + mocked_get_gpu_details, + mocked_env_cloud_details, + mocked_get_gpu_utilization_list, + mocked_is_gpu_details_available, + mocked_is_nvidia_system, + ): + # GIVEN a server that always fails to create a run, and a tracker + # whose cooldown will have already expired by the time stop() runs + # (no prior attempt at all here, so the very first attempt happens + # during stop()'s persist step). + responses.add( + responses.POST, + "http://test-api.com/runs", + json={"message": "down"}, + status=503, + ) + tracker = EmissionsTracker( + measure_power_secs=1, + save_to_file=False, + output_methods=[OutputMethod.API], + api_endpoint="http://test-api.com", + experiment_id="11111111-1111-1111-1111-111111111111", + api_key="fake-key", + ) + + # WHEN + tracker.start() + heavy_computation(run_time_secs=1) + tracker.stop() + + # THEN stop() should only attempt run creation once: the persist + # step's attempt already ran (and failed) this call, so the final + # bypass retry must be skipped instead of hitting the server again. + run_calls = [c for c in responses.calls if c.request.url.endswith("/runs")] + self.assertEqual(1, len(run_calls))