diff --git a/carbonserver/carbonserver/api/schemas.py b/carbonserver/carbonserver/api/schemas.py index a3343929d..74b267052 100644 --- a/carbonserver/carbonserver/api/schemas.py +++ b/carbonserver/carbonserver/api/schemas.py @@ -56,7 +56,7 @@ def __repr__(self): class EmissionBase(BaseModel): timestamp: datetime run_id: UUID - duration: int = Field( + duration: float = Field( ..., gt=0, description="The duration must be greater than zero" ) emissions_sum: Optional[float] = Field( @@ -283,7 +283,7 @@ class ExperimentReport(ExperimentBase): gpu_energy: float ram_energy: float energy_consumed: float - duration: int + duration: float emissions_rate: float emissions_count: int cpu_utilization_percent: Optional[float] = None @@ -396,7 +396,7 @@ class ProjectReport(ProjectBase): gpu_energy: float ram_energy: float energy_consumed: float - duration: int + duration: float emissions_rate: float emissions_count: int cpu_utilization_percent: Optional[float] = None @@ -443,7 +443,7 @@ class OrganizationReport(OrganizationBase): gpu_energy: float ram_energy: float energy_consumed: float - duration: int + duration: float emissions_rate: float emissions_count: int cpu_utilization_percent: Optional[float] = None diff --git a/carbonserver/tests/api/test_schema_compatibility.py b/carbonserver/tests/api/test_schema_compatibility.py index 94522fc8c..1acce76ca 100644 --- a/carbonserver/tests/api/test_schema_compatibility.py +++ b/carbonserver/tests/api/test_schema_compatibility.py @@ -118,3 +118,27 @@ def test_client_create_payloads_validate_against_server_schemas( client_payload, server_schema ): server_schema.model_validate(dataclasses.asdict(client_payload)) + + +def test_millisecond_duration_survives_client_to_server(): + """A sub-second duration is stored as is, not rounded away.""" + payload = client_schemas.EmissionCreate( + timestamp="2021-04-04T08:43:00+02:00", + run_id="40088f1a-d28e-4980-8d80-bf5600056a14", + duration=0.0042, + emissions_sum=1544.54, + emissions_rate=1.548444, + cpu_power=0.3, + gpu_power=0.0, + ram_power=0.15, + cpu_energy=55.21874, + gpu_energy=0.0, + ram_energy=2.0, + energy_consumed=57.21874, + ) + + validated = server_schemas.EmissionCreate.model_validate( + dataclasses.asdict(payload) + ) + + assert validated.duration == 0.0042 diff --git a/codecarbon/core/api_client.py b/codecarbon/core/api_client.py index bc2e0974e..3647e5e0f 100644 --- a/codecarbon/core/api_client.py +++ b/codecarbon/core/api_client.py @@ -189,15 +189,16 @@ def add_emission(self, carbon_emission: dict): "ApiClient.add_emission still no run_id, aborting for this time !" ) return False - if carbon_emission["duration"] < 1: + if carbon_emission["duration"] <= 0: + # The server declares duration as gt=0, so this would 422. logger.warning( - "ApiClient : emissions not sent because of a duration smaller than 1." + "ApiClient : emissions not sent because the duration is not positive." ) return False emission = EmissionCreate( timestamp=get_datetime_with_timezone(), run_id=self.run_id, - duration=int(carbon_emission["duration"]), + duration=carbon_emission["duration"], emissions_sum=carbon_emission["emissions"], emissions_rate=carbon_emission["emissions_rate"], cpu_power=carbon_emission["cpu_power"], @@ -215,7 +216,33 @@ 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) + response = requests.post( + url=url, json=payload, timeout=2, headers=self._get_headers() + ) + duration = payload["duration"] + if response.status_code == 422 and duration != round(duration): + # Servers older than the client declare duration as an int and + # reject fractional seconds. The old client only ever sent + # whole seconds, so reproduce that: below 1s there was never + # anything to send (rounding up would inflate a few-ms flush + # to a full second), otherwise retry once with the int value. + if duration < 1: + logger.debug( + "ApiClient : the API rejected a fractional duration" + " below 1s, it looks older than this client. Dropping" + " this emission instead of inflating its duration." + ) + return False + logger.info( + "ApiClient : the API rejected a fractional duration, it looks" + " older than this client. Retrying with the duration truncated to whole seconds." + ) + payload["duration"] = int(duration) + response = requests.post( + url=url, json=payload, timeout=2, headers=self._get_headers() + ) + if response.status_code != 201: + self._raise_api_error(url, payload, response) 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. diff --git a/codecarbon/core/schemas.py b/codecarbon/core/schemas.py index 84d1e9c77..9cc6b85a7 100644 --- a/codecarbon/core/schemas.py +++ b/codecarbon/core/schemas.py @@ -12,7 +12,7 @@ class EmissionBase: timestamp: str run_id: str - duration: int + duration: float emissions_sum: float emissions_rate: float cpu_power: float diff --git a/tests/test_api_call.py b/tests/test_api_call.py index 31e25c039..10adf056c 100644 --- a/tests/test_api_call.py +++ b/tests/test_api_call.py @@ -138,6 +138,68 @@ def test_call_api(self): assert payload["ram_utilization_percent"] == 56.5 assert payload["wue"] == 0.8 + def test_add_emission_retries_rounded_duration_on_old_server(self): + emission = { + "duration": 15.0023, + "emissions": 2.0, + "emissions_rate": 2.0, + "cpu_power": 3.0, + "gpu_power": 0, + "ram_power": 0.15, + "cpu_energy": 2, + "gpu_energy": 0, + "ram_energy": 1, + "energy_consumed": 3.0, + } + with requests_mock.Mocker() as m: + m.post( + "http://test.com/emissions", + [{"status_code": 422}, {"status_code": 201}], + ) + api = ApiClient( + endpoint_url="http://test.com", + experiment_id="exp-1", + create_run_automatically=False, + ) + api.run_id = "run-1" + + self.assertTrue(api.add_emission(emission)) + + self.assertEqual(m.call_count, 2) + self.assertEqual(m.request_history[0].json()["duration"], 15.0023) + retried = m.request_history[1].json()["duration"] + self.assertIsInstance(retried, int) + self.assertEqual(retried, 15) + + def test_add_emission_drops_subsecond_duration_on_old_server(self): + """An old server 422s on a fractional duration; below 1s we must drop + the emission (matching the old client, which never sent it), not + inflate it to 1s.""" + emission = { + "duration": 0.004, + "emissions": 2.0, + "emissions_rate": 2.0, + "cpu_power": 3.0, + "gpu_power": 0, + "ram_power": 0.15, + "cpu_energy": 2, + "gpu_energy": 0, + "ram_energy": 1, + "energy_consumed": 3.0, + } + with requests_mock.Mocker() as m: + m.post("http://test.com/emissions", status_code=422) + api = ApiClient( + endpoint_url="http://test.com", + experiment_id="exp-1", + create_run_automatically=False, + ) + api.run_id = "run-1" + + self.assertFalse(api.add_emission(emission)) + + self.assertEqual(m.call_count, 1) + def test_create_run_rounds_coordinates(self): with requests_mock.Mocker() as m: m.post("http://test.com/runs", json={"id": "run-1"}, status_code=201) @@ -235,31 +297,65 @@ def test_add_emission_returns_false_when_run_creation_fails(self): ) ) - def test_add_emission_skips_short_duration(self): - api = ApiClient( - endpoint_url="http://test.com", - experiment_id="exp-1", - conf=conf, - create_run_automatically=False, - ) - api.run_id = "run-1" + def test_add_emission_sends_millisecond_duration_unchanged(self): + """A sub-second duration is sent as is, not rounded.""" + with requests_mock.Mocker() as m: + m.post("http://test.com/emissions", json={"id": "em-1"}, 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" - self.assertFalse( - api.add_emission( - { - "duration": 0.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, - } + self.assertTrue( + api.add_emission( + { + "duration": 0.0042, + "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, + } + ) ) - ) + self.assertEqual(m.last_request.json()["duration"], 0.0042) + + def test_add_emission_skips_zero_duration(self): + """A zero-length flush would be a 422: the server requires duration > 0.""" + with requests_mock.Mocker() as m: + m.post("http://test.com/emissions", json={"id": "em-1"}, 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" + + self.assertFalse( + api.add_emission( + { + "duration": 0.0, + "emissions": 0.0, + "emissions_rate": 0.0, + "cpu_power": 1.0, + "gpu_power": 0.0, + "ram_power": 0.5, + "cpu_energy": 0.0, + "gpu_energy": 0.0, + "ram_energy": 0.0, + "energy_consumed": 0.0, + } + ) + ) + self.assertFalse(m.called) def test_add_emission_raises_on_unsuccessful_post(self): with requests_mock.Mocker() as m: