From bc75a1913acf5b0e758b2dfa4abf1fb54ced32cc Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rok=20Ro=C5=A1kar?= Date: Thu, 25 Jun 2026 17:53:54 +0200 Subject: [PATCH 1/7] poc: trying out openmeter metering --- .../renku_data_services/data_tasks/config.py | 21 ++++++ .../data_tasks/dependencies.py | 10 ++- .../resource_usage/core.py | 11 ++- .../resource_usage/metering.py | 74 +++++++++++++++++++ 4 files changed, 114 insertions(+), 2 deletions(-) create mode 100644 components/renku_data_services/resource_usage/metering.py diff --git a/bases/renku_data_services/data_tasks/config.py b/bases/renku_data_services/data_tasks/config.py index 7d0934065..ea405bbea 100644 --- a/bases/renku_data_services/data_tasks/config.py +++ b/bases/renku_data_services/data_tasks/config.py @@ -33,6 +33,24 @@ def from_env( return cls(enabled, api_key, host, environment) +@dataclass +class MeteringConfig: + """Configuration for the Kong metering endpoint.""" + + enabled: bool + endpoint_url: str + token: str + + @classmethod + def from_env(cls) -> "MeteringConfig": + """Create metering config from environment variables.""" + return cls( + enabled=os.environ.get("METERING_ENABLED", "false").lower() == "true", + endpoint_url=os.environ.get("METERING_ENDPOINT_URL", ""), + token=os.environ.get("METERING_API_TOKEN", ""), + ) + + @dataclass class Config: """Configuration for data tasks.""" @@ -40,6 +58,7 @@ class Config: db: DBConfig solr: SolrClientConfig posthog: PosthogConfig + metering: MeteringConfig authz: AuthzConfig keycloak: KeycloakConfig | None k8s_config_root: str @@ -67,6 +86,7 @@ def from_env(cls) -> Config: main_tick = int(os.environ.get("MAIN_LOG_INTERVAL_SECONDS", "300")) solr_config = SolrClientConfig.from_env() posthog_config = PosthogConfig.from_env() + metering_config = MeteringConfig.from_env() tcp_host = os.environ.get("TCP_HOST", "127.0.0.1") tcp_port = int(os.environ.get("TCP_PORT", "8001")) @@ -89,6 +109,7 @@ def from_env(cls) -> Config: main_log_interval_seconds=main_tick, solr=solr_config, posthog=posthog_config, + metering=metering_config, authz=authz, keycloak=keycloak, k8s_config_root=k8s_config_root, diff --git a/bases/renku_data_services/data_tasks/dependencies.py b/bases/renku_data_services/data_tasks/dependencies.py index c7a9ac036..04ce957d1 100644 --- a/bases/renku_data_services/data_tasks/dependencies.py +++ b/bases/renku_data_services/data_tasks/dependencies.py @@ -25,6 +25,7 @@ ResourcesRequestRecorder, ResourceUsageService, ) +from renku_data_services.resource_usage.metering import MeteringClient from renku_data_services.resource_usage.db import ResourceRequestsRepo from renku_data_services.search.db import SearchUpdatesRepo from renku_data_services.session.db import SessionRepository @@ -123,8 +124,15 @@ def from_env(cls, cfg: Config | None = None) -> "DependencyManager": resource_requests_recorder: ResourcesRequestRecorder if cfg.enable_resource_request_tracking: + metering_client = ( + MeteringClient(endpoint_url=cfg.metering.endpoint_url, token=cfg.metering.token) + if cfg.metering.enabled + else None + ) resource_requests_recorder = DefaultResourcesRequestRecorder( - repo=resource_requests_repo, fetch=ResourceRequestsFetch(k8s_client) + repo=resource_requests_repo, + fetch=ResourceRequestsFetch(k8s_client), + metering=metering_client, ) else: logger.warning("Resource request tracking is disabled!") diff --git a/components/renku_data_services/resource_usage/core.py b/components/renku_data_services/resource_usage/core.py index 91c106d39..f3470912e 100644 --- a/components/renku_data_services/resource_usage/core.py +++ b/components/renku_data_services/resource_usage/core.py @@ -11,6 +11,7 @@ from renku_data_services.k8s.models import GVK, K8sObject, K8sObjectFilter, K8sObjectMeta from renku_data_services.resource_usage import apispec from renku_data_services.resource_usage.db import ResourceRequestsRepo +from renku_data_services.resource_usage.metering import MeteringClient from renku_data_services.resource_usage.model import ( Credit, ResourceClassCost, @@ -155,9 +156,15 @@ async def record_resource_requests(self, interval: timedelta) -> None: class DefaultResourcesRequestRecorder(ResourcesRequestRecorder): """Methods for recording resource requests.""" - def __init__(self, repo: ResourceRequestsRepo, fetch: ResourceRequestsFetchProto) -> None: + def __init__( + self, + repo: ResourceRequestsRepo, + fetch: ResourceRequestsFetchProto, + metering: MeteringClient | None = None, + ) -> None: self._repo = repo self._fetch = fetch + self._metering = metering async def record_resource_requests(self, interval: timedelta) -> None: """Fetches all resource requests in the given namespace and stores them.""" @@ -170,6 +177,8 @@ async def record_resource_requests(self, interval: timedelta) -> None: else: logger.info(f"Inserting {size} resource request records.") await self._repo.insert_many(result) + if self._metering is not None: + await self._metering.emit(result) class ResourceUsageService: diff --git a/components/renku_data_services/resource_usage/metering.py b/components/renku_data_services/resource_usage/metering.py new file mode 100644 index 000000000..82c71211a --- /dev/null +++ b/components/renku_data_services/resource_usage/metering.py @@ -0,0 +1,74 @@ +"""CloudEvent emission for session resource usage metering.""" + +from datetime import UTC, datetime + +import httpx + +from renku_data_services.app_config import logging +from renku_data_services.resource_usage.model import ResourcesRequest + +logger = logging.getLogger(__file__) + +_CLOUDEVENTS_BATCH_CONTENT_TYPE = "application/cloudevents-batch+json" +_CLOUDEVENTS_DATA_CONTENT_TYPE = "application/json" +_CLOUDEVENTS_SPEC_VERSION = "1.0" +_CLOUDEVENTS_SOURCE = "/renku-data-services/resource-usage" + + +def _to_cloudevent(req: ResourcesRequest) -> dict: + event_id = f"{req.uid}/{req.capture_date.isoformat()}" + data: dict = { + "uid": req.uid, + "kind": req.kind, + "user_id": req.user_id, + "project_id": str(req.project_id) if req.project_id is not None else None, + "launcher_id": str(req.launcher_id) if req.launcher_id is not None else None, + "resource_class_id": req.resource_class_id, + "resource_pool_id": req.resource_pool_id, + "cluster_id": str(req.cluster_id) if req.cluster_id is not None else None, + "phase": req.phase, + "capture_interval_seconds": req.capture_interval.total_seconds(), + "cpu_millicores": req.data.cpu.milli_cores if req.data.cpu is not None else None, + "memory_bytes": req.data.memory.bytes if req.data.memory is not None else None, + "gpu_cores": req.data.gpu.cores if req.data.gpu is not None else None, + "disk_bytes": req.data.disk.bytes if req.data.disk is not None else None, + } + return { + "specversion": _CLOUDEVENTS_SPEC_VERSION, + "type": "ch.renku.session.resource_usage", + "source": _CLOUDEVENTS_SOURCE, + "id": event_id, + "time": req.capture_date.astimezone(UTC).isoformat(), + "datacontenttype": _CLOUDEVENTS_DATA_CONTENT_TYPE, + "data": data, + } + + +class MeteringClient: + """Emits session resource usage as CloudEvents to a metering endpoint.""" + + def __init__(self, endpoint_url: str, token: str) -> None: + self._endpoint_url = endpoint_url + self._headers = { + "Authorization": f"Bearer {token}", + "Content-Type": _CLOUDEVENTS_BATCH_CONTENT_TYPE, + } + + async def emit(self, requests: list[ResourcesRequest]) -> None: + """POST all resource requests as a CloudEvents batch. Never raises.""" + if not requests: + return + events = [_to_cloudevent(r) for r in requests] + try: + async with httpx.AsyncClient(timeout=30) as client: + resp = await client.post(self._endpoint_url, headers=self._headers, json=events) + if resp.status_code >= 300 or resp.status_code < 200: + logger.warning( + f"Metering endpoint returned unexpected status {resp.status_code}: {resp.text[:200]}" + ) + else: + logger.debug(f"Emitted {len(events)} metering events, status={resp.status_code}") + except httpx.HTTPError as ex: + logger.warning(f"Failed to emit metering events: {ex}", exc_info=ex) + except Exception as ex: + logger.warning(f"Unexpected error emitting metering events: {ex}", exc_info=ex) From 6b5db52f82c9892ed3a2380ac427aa0e239a2702 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rok=20Ro=C5=A1kar?= Date: Thu, 25 Jun 2026 19:03:04 +0200 Subject: [PATCH 2/7] update event logic --- .../resource_usage/core.py | 4 +- .../renku_data_services/resource_usage/db.py | 5 ++ .../resource_usage/metering.py | 55 ++++++++++--------- 3 files changed, 37 insertions(+), 27 deletions(-) diff --git a/components/renku_data_services/resource_usage/core.py b/components/renku_data_services/resource_usage/core.py index f3470912e..8abe8f1ea 100644 --- a/components/renku_data_services/resource_usage/core.py +++ b/components/renku_data_services/resource_usage/core.py @@ -178,7 +178,9 @@ async def record_resource_requests(self, interval: timedelta) -> None: logger.info(f"Inserting {size} resource request records.") await self._repo.insert_many(result) if self._metering is not None: - await self._metering.emit(result) + class_ids = {r.resource_class_id for r in result if r.resource_class_id is not None} + costs = await self._repo.get_costs_by_class_ids(class_ids) + await self._metering.emit(result, costs) class ResourceUsageService: diff --git a/components/renku_data_services/resource_usage/db.py b/components/renku_data_services/resource_usage/db.py index d402f2ac8..ac745e646 100644 --- a/components/renku_data_services/resource_usage/db.py +++ b/components/renku_data_services/resource_usage/db.py @@ -84,6 +84,11 @@ async def insert_many(self, reqs: Iterable[ResourcesRequest]) -> None: session.add_all(vals) await session.flush() + async def get_costs_by_class_ids(self, ids: set[int]) -> dict[int, Credit]: + """Return a mapping of resource_class_id to Credit for the given ids.""" + async with self.session_maker() as session: + return await self._get_all_costs(session, ids) + async def _get_all_costs(self, session: AsyncSession, ids: set[int]) -> dict[int, Credit]: stmt = sa.select(ResourceClassCostORM.id, ResourceClassCostORM.cost).where(ResourceClassCostORM.id.in_(ids)) rows = await session.execute(stmt) diff --git a/components/renku_data_services/resource_usage/metering.py b/components/renku_data_services/resource_usage/metering.py index 82c71211a..a25d0ad74 100644 --- a/components/renku_data_services/resource_usage/metering.py +++ b/components/renku_data_services/resource_usage/metering.py @@ -1,11 +1,11 @@ """CloudEvent emission for session resource usage metering.""" -from datetime import UTC, datetime +from datetime import UTC import httpx from renku_data_services.app_config import logging -from renku_data_services.resource_usage.model import ResourcesRequest +from renku_data_services.resource_usage.model import Credit, ResourcesRequest logger = logging.getLogger(__file__) @@ -13,35 +13,38 @@ _CLOUDEVENTS_DATA_CONTENT_TYPE = "application/json" _CLOUDEVENTS_SPEC_VERSION = "1.0" _CLOUDEVENTS_SOURCE = "/renku-data-services/resource-usage" +_CLOUDEVENTS_TYPE = "ch.renku.session.resource_usage" -def _to_cloudevent(req: ResourcesRequest) -> dict: - event_id = f"{req.uid}/{req.capture_date.isoformat()}" - data: dict = { - "uid": req.uid, - "kind": req.kind, - "user_id": req.user_id, - "project_id": str(req.project_id) if req.project_id is not None else None, - "launcher_id": str(req.launcher_id) if req.launcher_id is not None else None, - "resource_class_id": req.resource_class_id, - "resource_pool_id": req.resource_pool_id, - "cluster_id": str(req.cluster_id) if req.cluster_id is not None else None, - "phase": req.phase, - "capture_interval_seconds": req.capture_interval.total_seconds(), - "cpu_millicores": req.data.cpu.milli_cores if req.data.cpu is not None else None, - "memory_bytes": req.data.memory.bytes if req.data.memory is not None else None, - "gpu_cores": req.data.gpu.cores if req.data.gpu is not None else None, - "disk_bytes": req.data.disk.bytes if req.data.disk is not None else None, - } - return { +def _to_cloudevent(req: ResourcesRequest, costs: dict[int, Credit]) -> dict: + cu_cost: float | None = None + if req.resource_class_id is not None: + cost = costs.get(req.resource_class_id, Credit.zero()) + cu_cost = round(cost.value * (req.capture_interval.total_seconds() / 3600.0), 6) + + event: dict = { "specversion": _CLOUDEVENTS_SPEC_VERSION, - "type": "ch.renku.session.resource_usage", + "type": _CLOUDEVENTS_TYPE, "source": _CLOUDEVENTS_SOURCE, - "id": event_id, + "id": f"{req.uid}/{req.capture_date.astimezone(UTC).isoformat()}", "time": req.capture_date.astimezone(UTC).isoformat(), "datacontenttype": _CLOUDEVENTS_DATA_CONTENT_TYPE, - "data": data, + "data": { + "uid": req.uid, + "kind": req.kind, + "project_id": str(req.project_id) if req.project_id is not None else None, + "launcher_id": str(req.launcher_id) if req.launcher_id is not None else None, + "resource_class_id": req.resource_class_id, + "resource_pool_id": req.resource_pool_id, + "cluster_id": str(req.cluster_id) if req.cluster_id is not None else None, + "phase": req.phase, + "capture_interval_seconds": req.capture_interval.total_seconds(), + "cu_cost": cu_cost, + }, } + if req.user_id is not None: + event["subject"] = req.user_id + return event class MeteringClient: @@ -54,11 +57,11 @@ def __init__(self, endpoint_url: str, token: str) -> None: "Content-Type": _CLOUDEVENTS_BATCH_CONTENT_TYPE, } - async def emit(self, requests: list[ResourcesRequest]) -> None: + async def emit(self, requests: list[ResourcesRequest], costs: dict[int, Credit]) -> None: """POST all resource requests as a CloudEvents batch. Never raises.""" if not requests: return - events = [_to_cloudevent(r) for r in requests] + events = [_to_cloudevent(r, costs) for r in requests] try: async with httpx.AsyncClient(timeout=30) as client: resp = await client.post(self._endpoint_url, headers=self._headers, json=events) From d418c6bb02883cf809490983b040f7456f93ad47 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rok=20Ro=C5=A1kar?= Date: Fri, 26 Jun 2026 11:03:23 +0200 Subject: [PATCH 3/7] fix: skip metering events without a resource class id Pods/PVCs without a resource_class_id produce a null cu_cost which Kong Konnect rejects with a 400. Filter these out before emitting. --- components/renku_data_services/resource_usage/metering.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/components/renku_data_services/resource_usage/metering.py b/components/renku_data_services/resource_usage/metering.py index a25d0ad74..0c1e5dfee 100644 --- a/components/renku_data_services/resource_usage/metering.py +++ b/components/renku_data_services/resource_usage/metering.py @@ -61,7 +61,7 @@ async def emit(self, requests: list[ResourcesRequest], costs: dict[int, Credit]) """POST all resource requests as a CloudEvents batch. Never raises.""" if not requests: return - events = [_to_cloudevent(r, costs) for r in requests] + events = [_to_cloudevent(r, costs) for r in requests if r.resource_class_id is not None] try: async with httpx.AsyncClient(timeout=30) as client: resp = await client.post(self._endpoint_url, headers=self._headers, json=events) From 1d4989dced29f13868ce8dade535b7231797d59e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rok=20Ro=C5=A1kar?= Date: Fri, 26 Jun 2026 13:31:52 +0200 Subject: [PATCH 4/7] chore: increase logging and reduce interval --- bases/renku_data_services/data_tasks/task_defs.py | 2 +- components/renku_data_services/resource_usage/metering.py | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/bases/renku_data_services/data_tasks/task_defs.py b/bases/renku_data_services/data_tasks/task_defs.py index 8f080e254..f1b7116dd 100644 --- a/bases/renku_data_services/data_tasks/task_defs.py +++ b/bases/renku_data_services/data_tasks/task_defs.py @@ -440,7 +440,7 @@ async def cleanup_orphaned_capacity_reservations(dm: DependencyManager) -> None: async def record_resource_requests(dm: DependencyManager) -> None: """Periodically record all resource requests.""" - interval_seconds = 600 + interval_seconds = 30 # increase before merging while True: await dm.resource_requests_recorder.record_resource_requests(timedelta(seconds=interval_seconds)) await asyncio.sleep(interval_seconds) diff --git a/components/renku_data_services/resource_usage/metering.py b/components/renku_data_services/resource_usage/metering.py index 0c1e5dfee..8e815d590 100644 --- a/components/renku_data_services/resource_usage/metering.py +++ b/components/renku_data_services/resource_usage/metering.py @@ -70,7 +70,7 @@ async def emit(self, requests: list[ResourcesRequest], costs: dict[int, Credit]) f"Metering endpoint returned unexpected status {resp.status_code}: {resp.text[:200]}" ) else: - logger.debug(f"Emitted {len(events)} metering events, status={resp.status_code}") + logger.info(f"Emitted {len(events)} metering events, status={resp.status_code}") except httpx.HTTPError as ex: logger.warning(f"Failed to emit metering events: {ex}", exc_info=ex) except Exception as ex: From ce76db0cd8eb5825b57c8bdc475cb95170560bea Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rok=20Ro=C5=A1kar?= Date: Wed, 1 Jul 2026 14:14:47 +0200 Subject: [PATCH 5/7] feat: send events to meteroid --- .../resource_usage/core.py | 4 +- .../resource_usage/metering.py | 86 ++++++++++--------- 2 files changed, 47 insertions(+), 43 deletions(-) diff --git a/components/renku_data_services/resource_usage/core.py b/components/renku_data_services/resource_usage/core.py index 8abe8f1ea..c4dc31ec1 100644 --- a/components/renku_data_services/resource_usage/core.py +++ b/components/renku_data_services/resource_usage/core.py @@ -11,7 +11,7 @@ from renku_data_services.k8s.models import GVK, K8sObject, K8sObjectFilter, K8sObjectMeta from renku_data_services.resource_usage import apispec from renku_data_services.resource_usage.db import ResourceRequestsRepo -from renku_data_services.resource_usage.metering import MeteringClient +from renku_data_services.resource_usage.metering import MeteringClient, MetricCode from renku_data_services.resource_usage.model import ( Credit, ResourceClassCost, @@ -180,7 +180,7 @@ async def record_resource_requests(self, interval: timedelta) -> None: if self._metering is not None: class_ids = {r.resource_class_id for r in result if r.resource_class_id is not None} costs = await self._repo.get_costs_by_class_ids(class_ids) - await self._metering.emit(result, costs) + await self._metering.emit(result, costs, MetricCode.session_resource_usage) class ResourceUsageService: diff --git a/components/renku_data_services/resource_usage/metering.py b/components/renku_data_services/resource_usage/metering.py index 8e815d590..df495b69f 100644 --- a/components/renku_data_services/resource_usage/metering.py +++ b/components/renku_data_services/resource_usage/metering.py @@ -1,6 +1,7 @@ -"""CloudEvent emission for session resource usage metering.""" +"""Meteroid event emission for session resource usage metering.""" from datetime import UTC +from enum import StrEnum import httpx @@ -9,62 +10,65 @@ logger = logging.getLogger(__file__) -_CLOUDEVENTS_BATCH_CONTENT_TYPE = "application/cloudevents-batch+json" -_CLOUDEVENTS_DATA_CONTENT_TYPE = "application/json" -_CLOUDEVENTS_SPEC_VERSION = "1.0" -_CLOUDEVENTS_SOURCE = "/renku-data-services/resource-usage" -_CLOUDEVENTS_TYPE = "ch.renku.session.resource_usage" +_INGEST_PATH = "/api/v1/events/ingest" -def _to_cloudevent(req: ResourcesRequest, costs: dict[int, Credit]) -> dict: - cu_cost: float | None = None - if req.resource_class_id is not None: - cost = costs.get(req.resource_class_id, Credit.zero()) - cu_cost = round(cost.value * (req.capture_interval.total_seconds() / 3600.0), 6) +class MetricCode(StrEnum): + session_resource_usage = "session_resource_usage" - event: dict = { - "specversion": _CLOUDEVENTS_SPEC_VERSION, - "type": _CLOUDEVENTS_TYPE, - "source": _CLOUDEVENTS_SOURCE, - "id": f"{req.uid}/{req.capture_date.astimezone(UTC).isoformat()}", - "time": req.capture_date.astimezone(UTC).isoformat(), - "datacontenttype": _CLOUDEVENTS_DATA_CONTENT_TYPE, - "data": { - "uid": req.uid, - "kind": req.kind, - "project_id": str(req.project_id) if req.project_id is not None else None, - "launcher_id": str(req.launcher_id) if req.launcher_id is not None else None, - "resource_class_id": req.resource_class_id, - "resource_pool_id": req.resource_pool_id, - "cluster_id": str(req.cluster_id) if req.cluster_id is not None else None, - "phase": req.phase, - "capture_interval_seconds": req.capture_interval.total_seconds(), - "cu_cost": cu_cost, - }, + +def _to_meteroid_event(req: ResourcesRequest, costs: dict[int, Credit], metric_code: str) -> dict: + cost = costs.get(req.resource_class_id, Credit.zero()) # type: ignore[arg-type] + cu_cost = round(cost.value * (req.capture_interval.total_seconds() / 3600.0), 6) + + properties: dict[str, str] = { + "cu_cost": str(cu_cost), + "kind": req.kind, + "phase": req.phase, + "capture_interval_seconds": str(req.capture_interval.total_seconds()), + "resource_class_id": str(req.resource_class_id), + } + if req.resource_pool_id is not None: + properties["resource_pool_id"] = str(req.resource_pool_id) + if req.project_id is not None: + properties["project_id"] = str(req.project_id) + if req.launcher_id is not None: + properties["launcher_id"] = str(req.launcher_id) + if req.cluster_id is not None: + properties["cluster_id"] = str(req.cluster_id) + + return { + "event_id": f"{req.uid}/{req.capture_date.astimezone(UTC).isoformat()}", + "code": metric_code, + "customer_id": req.user_id, + "timestamp": req.capture_date.astimezone(UTC).isoformat(), + "properties": properties, } - if req.user_id is not None: - event["subject"] = req.user_id - return event class MeteringClient: - """Emits session resource usage as CloudEvents to a metering endpoint.""" + """Emits session resource usage events to Meteroid.""" def __init__(self, endpoint_url: str, token: str) -> None: - self._endpoint_url = endpoint_url + self._endpoint_url = endpoint_url.rstrip("/") + _INGEST_PATH self._headers = { "Authorization": f"Bearer {token}", - "Content-Type": _CLOUDEVENTS_BATCH_CONTENT_TYPE, + "Content-Type": "application/json", } - async def emit(self, requests: list[ResourcesRequest], costs: dict[int, Credit]) -> None: - """POST all resource requests as a CloudEvents batch. Never raises.""" - if not requests: + async def emit(self, requests: list[ResourcesRequest], costs: dict[int, Credit], metric_code: MetricCode) -> None: + """POST all resource requests as a Meteroid ingest batch. Never raises.""" + events = [ + _to_meteroid_event(r, costs, metric_code) + for r in requests + if r.resource_class_id is not None and r.user_id is not None + ] + if not events: return - events = [_to_cloudevent(r, costs) for r in requests if r.resource_class_id is not None] + body = {"allow_partial_failures": True, "events": events} try: async with httpx.AsyncClient(timeout=30) as client: - resp = await client.post(self._endpoint_url, headers=self._headers, json=events) + resp = await client.post(self._endpoint_url, headers=self._headers, json=body) if resp.status_code >= 300 or resp.status_code < 200: logger.warning( f"Metering endpoint returned unexpected status {resp.status_code}: {resp.text[:200]}" From e44c763829a1f88b62829bca510e9a6c78482e99 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rok=20Ro=C5=A1kar?= Date: Wed, 1 Jul 2026 14:32:58 +0200 Subject: [PATCH 6/7] chore: fine tune events and print json body --- components/renku_data_services/resource_usage/metering.py | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/components/renku_data_services/resource_usage/metering.py b/components/renku_data_services/resource_usage/metering.py index df495b69f..d0565edab 100644 --- a/components/renku_data_services/resource_usage/metering.py +++ b/components/renku_data_services/resource_usage/metering.py @@ -10,9 +10,6 @@ logger = logging.getLogger(__file__) -_INGEST_PATH = "/api/v1/events/ingest" - - class MetricCode(StrEnum): session_resource_usage = "session_resource_usage" @@ -50,7 +47,7 @@ class MeteringClient: """Emits session resource usage events to Meteroid.""" def __init__(self, endpoint_url: str, token: str) -> None: - self._endpoint_url = endpoint_url.rstrip("/") + _INGEST_PATH + self._endpoint_url = endpoint_url self._headers = { "Authorization": f"Bearer {token}", "Content-Type": "application/json", @@ -74,7 +71,7 @@ async def emit(self, requests: list[ResourcesRequest], costs: dict[int, Credit], f"Metering endpoint returned unexpected status {resp.status_code}: {resp.text[:200]}" ) else: - logger.info(f"Emitted {len(events)} metering events, status={resp.status_code}") + logger.info(f"Emitted {len(events)} metering events, status={resp.status_code}: {body}") except httpx.HTTPError as ex: logger.warning(f"Failed to emit metering events: {ex}", exc_info=ex) except Exception as ex: From 7b9dcc619e75d1b9b308ddffdf2b665c86c3c034 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rok=20Ro=C5=A1kar?= Date: Wed, 1 Jul 2026 15:02:54 +0200 Subject: [PATCH 7/7] chore: use resource-pool-id as a customer-id proxy for now --- components/renku_data_services/resource_usage/metering.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/components/renku_data_services/resource_usage/metering.py b/components/renku_data_services/resource_usage/metering.py index d0565edab..c0a7ae417 100644 --- a/components/renku_data_services/resource_usage/metering.py +++ b/components/renku_data_services/resource_usage/metering.py @@ -24,6 +24,7 @@ def _to_meteroid_event(req: ResourcesRequest, costs: dict[int, Credit], metric_c "phase": req.phase, "capture_interval_seconds": str(req.capture_interval.total_seconds()), "resource_class_id": str(req.resource_class_id), + "user_id": str(req.user_id), } if req.resource_pool_id is not None: properties["resource_pool_id"] = str(req.resource_pool_id) @@ -37,7 +38,7 @@ def _to_meteroid_event(req: ResourcesRequest, costs: dict[int, Credit], metric_c return { "event_id": f"{req.uid}/{req.capture_date.astimezone(UTC).isoformat()}", "code": metric_code, - "customer_id": req.user_id, + "customer_id": f"resource_pool_id-{req.resource_pool_id}", "timestamp": req.capture_date.astimezone(UTC).isoformat(), "properties": properties, } @@ -58,7 +59,7 @@ async def emit(self, requests: list[ResourcesRequest], costs: dict[int, Credit], events = [ _to_meteroid_event(r, costs, metric_code) for r in requests - if r.resource_class_id is not None and r.user_id is not None + if r.resource_class_id is not None and r.user_id is not None and r.resource_pool_id is not None ] if not events: return