diff --git a/codecarbon/core/api_client.py b/codecarbon/core/api_client.py index bc2e0974e..9b89d7527 100644 --- a/codecarbon/core/api_client.py +++ b/codecarbon/core/api_client.py @@ -8,7 +8,8 @@ # from httpx import AsyncClient import dataclasses import json -from datetime import timedelta, tzinfo +import time +from datetime import datetime, timedelta, tzinfo import requests @@ -33,6 +34,29 @@ 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) +# 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: + """ + 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,11 +82,14 @@ 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 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) @@ -80,16 +107,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 +133,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 +165,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 +173,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 +181,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 +197,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,28 +205,44 @@ 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): + 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." ) 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 +262,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. @@ -225,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 ? @@ -235,6 +286,15 @@ def _create_run(self, experiment_id: str): "ApiClient FATAL The ApiClient._create_run() needs an experiment_id !" ) return None + if ( + 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") + 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(), @@ -256,8 +316,11 @@ 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"] + 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" @@ -282,7 +345,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 +361,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 +369,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/emissions_tracker.py b/codecarbon/emissions_tracker.py index 96ed00c91..de638dacb 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 @@ -925,12 +926,28 @@ 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, experiment_name=self._experiment_name, ) + # If run creation was still in its cooldown, the emission above was + # 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: + 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/codecarbon/output_methods/http.py b/codecarbon/output_methods/http.py index e0ff710b1..bc2866641 100644 --- a/codecarbon/output_methods/http.py +++ b/codecarbon/output_methods/http.py @@ -57,20 +57,23 @@ def __init__( ) self.run_id = self.api.run_id - def _ensure_api_run(self) -> None: + def exit(self) -> None: + self.api.close() + + 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 31e25c039..ac042128f 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 @@ -216,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( @@ -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) @@ -303,6 +343,133 @@ 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_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) 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() 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))