396 lines
16 KiB
Python
396 lines
16 KiB
Python
"""
|
|
providers.py — Accès aux API de coûts AWS et Azure, + conversion de devise.
|
|
|
|
Module partagé entre le collecteur cron (collect.py), l'import historique
|
|
(import_historical.py) et les scripts CLI ponctuels. Toute la logique
|
|
d'authentification et d'agrégation vit ici, une seule fois.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
|
|
class ResourceLevelUnavailable(Exception):
|
|
"""Levée quand AWS/Azure refuse une demande de détail par ressource —
|
|
le plus souvent parce que "Resource IDs" n'est pas activé côté AWS
|
|
Cost Explorer, ou qu'une fenêtre de plus de 14 jours a été demandée à
|
|
AWS. Distincte d'une erreur générique pour que l'appelant puisse
|
|
l'afficher comme un état "indisponible", pas comme un échec."""
|
|
|
|
|
|
def _short_resource_name(resource_id: str) -> str:
|
|
"""Dernier segment d'un ARN AWS ou d'un resource ID Azure, pour
|
|
l'affichage (ex: 'arn:aws:s3:::my-bucket' -> 'my-bucket',
|
|
'/subscriptions/.../virtualMachines/vm-1' -> 'vm-1')."""
|
|
if "/" in resource_id:
|
|
return resource_id.rsplit("/", 1)[-1]
|
|
if ":" in resource_id:
|
|
return resource_id.rsplit(":", 1)[-1]
|
|
return resource_id
|
|
|
|
|
|
# ----------------------------------------------------------------------
|
|
# AWS
|
|
# ----------------------------------------------------------------------
|
|
|
|
def _aws_session(entry: dict):
|
|
"""Construit la session boto3 pour un tenant AWS.
|
|
|
|
Deux modes, au choix par tenant dans config.json :
|
|
- 'access_key_id' + 'secret_access_key' : clé IAM dédiée au tenant
|
|
(droit ce:GetCostAndUsage uniquement) — le seul mode qui fonctionne
|
|
en conteneur/Kubernetes, car il ne dépend d'aucun état sur l'hôte.
|
|
- 'profile' : profil ~/.aws/{credentials,config} existant — pratique en
|
|
local, mais suppose ~/.aws monté dans le conteneur, et ne fonctionne
|
|
pas du tout avec un profil SSO (le refresh token expire et sa
|
|
reconduction demande un navigateur, indisponible en conteneur).
|
|
"""
|
|
import boto3
|
|
|
|
name = entry.get("name", "?")
|
|
if entry.get("access_key_id") and entry.get("secret_access_key"):
|
|
return boto3.Session(
|
|
aws_access_key_id=entry["access_key_id"],
|
|
aws_secret_access_key=entry["secret_access_key"],
|
|
)
|
|
if entry.get("profile"):
|
|
return boto3.Session(profile_name=entry["profile"])
|
|
raise ValueError(
|
|
f"Tenant AWS '{name}' : config.json doit fournir soit 'profile', "
|
|
"soit 'access_key_id' + 'secret_access_key'."
|
|
)
|
|
|
|
|
|
def get_aws_cost(entry: dict, start: str, end: str, granularity: str = "MONTHLY") -> dict:
|
|
"""Coût total + répartition par service pour un tenant AWS (voir
|
|
_aws_session pour les modes d'authentification), sur la période
|
|
[start, end) (end exclusif, format YYYY-MM-DD).
|
|
|
|
granularity: 'DAILY' ou 'MONTHLY' (voir doc AWS Cost Explorer).
|
|
"""
|
|
session = _aws_session(entry)
|
|
ce = session.client("ce", region_name="us-east-1") # Cost Explorer = toujours us-east-1
|
|
|
|
resp = ce.get_cost_and_usage(
|
|
TimePeriod={"Start": start, "End": end},
|
|
Granularity=granularity,
|
|
Metrics=["UnblendedCost"],
|
|
GroupBy=[{"Type": "DIMENSION", "Key": "SERVICE"}],
|
|
)
|
|
|
|
total = 0.0
|
|
by_service: dict[str, float] = {}
|
|
currency = "USD"
|
|
for period in resp["ResultsByTime"]:
|
|
for group in period["Groups"]:
|
|
service = group["Keys"][0]
|
|
amount = float(group["Metrics"]["UnblendedCost"]["Amount"])
|
|
currency = group["Metrics"]["UnblendedCost"]["Unit"]
|
|
by_service[service] = by_service.get(service, 0.0) + amount
|
|
total += amount
|
|
|
|
return {
|
|
"total": round(total, 2),
|
|
"currency": currency,
|
|
"by_service": {k: round(v, 2) for k, v in sorted(by_service.items(), key=lambda x: -x[1])},
|
|
}
|
|
|
|
|
|
def get_aws_resource_costs(entry: dict, start: str, end: str) -> dict:
|
|
"""Coût par ressource (bucket S3, instance EC2, fonction Lambda...),
|
|
groupé par service, sur la période [start, end). Le détail par ressource
|
|
n'est PAS accessible via get_cost_and_usage (son GroupBy ne connaît pas
|
|
RESOURCE_ID) — il faut l'API dédiée get_cost_and_usage_with_resources,
|
|
qui impose les mêmes contraintes (14 jours max, "Resource IDs" activé
|
|
côté compte). Lève ResourceLevelUnavailable si AWS refuse, plutôt que de
|
|
laisser planter l'appelant, pour que ce soit affichable comme un état
|
|
plutôt qu'une erreur de collecte."""
|
|
from botocore.exceptions import ClientError
|
|
|
|
session = _aws_session(entry)
|
|
ce = session.client("ce", region_name="us-east-1")
|
|
|
|
try:
|
|
resp = ce.get_cost_and_usage_with_resources(
|
|
TimePeriod={"Start": start, "End": end},
|
|
Granularity="DAILY",
|
|
Metrics=["UnblendedCost"],
|
|
GroupBy=[
|
|
{"Type": "DIMENSION", "Key": "SERVICE"},
|
|
{"Type": "DIMENSION", "Key": "RESOURCE_ID"},
|
|
],
|
|
Filter={"Dimensions": {"Key": "RECORD_TYPE", "Values": ["Usage"]}},
|
|
)
|
|
except ClientError as e:
|
|
raise ResourceLevelUnavailable(str(e)) from e
|
|
|
|
services: dict[str, dict[str, float]] = {}
|
|
currency = "USD"
|
|
for period in resp["ResultsByTime"]:
|
|
for group in period["Groups"]:
|
|
service, resource_id = group["Keys"]
|
|
amount = float(group["Metrics"]["UnblendedCost"]["Amount"])
|
|
currency = group["Metrics"]["UnblendedCost"]["Unit"]
|
|
bucket = services.setdefault(service, {})
|
|
bucket[resource_id] = bucket.get(resource_id, 0.0) + amount
|
|
|
|
return {"currency": currency, "services": _finalize_resource_buckets(services)}
|
|
|
|
|
|
def _finalize_resource_buckets(services: dict[str, dict[str, float]]) -> dict[str, list[dict]]:
|
|
"""Transforme {service: {resource_id: montant}} en
|
|
{service: [{resource_id, resource_name, amount}, ...]} trié par coût
|
|
décroissant. 'NoResourceId' / id vide = coûts du service non rattachés
|
|
à une ressource précise (data transfer, support...)."""
|
|
result: dict[str, list[dict]] = {}
|
|
for service, resources in services.items():
|
|
rows = []
|
|
for rid, amount in resources.items():
|
|
is_unattributed = not rid or rid == "NoResourceId"
|
|
rows.append({
|
|
"resource_id": rid,
|
|
"resource_name": "Autres coûts (sans ressource associée)" if is_unattributed else _short_resource_name(rid),
|
|
"amount": round(amount, 2),
|
|
})
|
|
rows.sort(key=lambda r: -r["amount"])
|
|
result[service] = rows
|
|
return result
|
|
|
|
|
|
# ----------------------------------------------------------------------
|
|
# Azure
|
|
# ----------------------------------------------------------------------
|
|
|
|
# Cost Management ne renvoie pas le `Retry-After` HTTP standard sur ses 429 —
|
|
# constaté en pratique : il renvoie ses propres en-têtes `x-ms-ratelimit-
|
|
# microsoft.costmanagement-*-retry-after`, un par niveau de quota (entité,
|
|
# tenant, type de client). On vérifie les deux formes, et on prend le délai
|
|
# le plus grand si plusieurs sont présents (le plus restrictif fait foi).
|
|
_AZURE_RETRY_AFTER_HEADERS = (
|
|
"Retry-After",
|
|
"x-ms-ratelimit-microsoft.costmanagement-entity-retry-after",
|
|
"x-ms-ratelimit-microsoft.costmanagement-tenant-retry-after",
|
|
"x-ms-ratelimit-microsoft.costmanagement-clienttype-retry-after",
|
|
)
|
|
|
|
|
|
def _azure_retry_after_seconds(response) -> float | None:
|
|
if response is None:
|
|
return None
|
|
waits = []
|
|
for header in _AZURE_RETRY_AFTER_HEADERS:
|
|
value = response.headers.get(header)
|
|
if value:
|
|
try:
|
|
waits.append(float(value))
|
|
except ValueError:
|
|
pass
|
|
return max(waits) if waits else None
|
|
|
|
|
|
def _azure_query_with_retry(client, scope: str, query: dict, max_retries: int = 8):
|
|
"""client.query.usage avec retry + backoff sur les 429. L'API Cost
|
|
Management est plus sujette au rate-limiting que Cost Explorer côté AWS
|
|
— le quota est souvent partagé au niveau du tenant Azure AD entier, pas
|
|
par service principal, donc même un tout premier appel peut être
|
|
limité — et azure-core ne retente pas toujours ces réponses tout seul
|
|
(constaté : un seul essai visible dans les logs avant l'erreur).
|
|
|
|
max_retries=8 avec un plafond de repli à 60s (~5 min de patience max
|
|
cumulée) si Azure ne fournit aucun en-tête retry-after exploitable :
|
|
volontairement généreux, pensé pour un backfill qui peut se permettre
|
|
d'attendre plutôt que d'abandonner en cours de route."""
|
|
from azure.core.exceptions import HttpResponseError
|
|
import time
|
|
|
|
for attempt in range(max_retries):
|
|
try:
|
|
return client.query.usage(scope, query)
|
|
except HttpResponseError as e:
|
|
if e.status_code != 429 or attempt == max_retries - 1:
|
|
raise
|
|
wait = _azure_retry_after_seconds(e.response)
|
|
if wait is None:
|
|
wait = min(2 ** attempt, 60)
|
|
time.sleep(wait)
|
|
|
|
|
|
def get_azure_cost(tenant_id: str, client_id: str, client_secret: str,
|
|
subscription_id: str, start: str, end: str) -> dict:
|
|
"""Coût total + répartition par service pour une subscription Azure,
|
|
sur la période [start, end)."""
|
|
from azure.identity import ClientSecretCredential
|
|
from azure.mgmt.costmanagement import CostManagementClient
|
|
|
|
credential = ClientSecretCredential(tenant_id, client_id, client_secret)
|
|
client = CostManagementClient(credential)
|
|
scope = f"/subscriptions/{subscription_id}"
|
|
|
|
query = {
|
|
"type": "ActualCost",
|
|
"timeframe": "Custom",
|
|
"timePeriod": {"from": f"{start}T00:00:00+00:00", "to": f"{end}T00:00:00+00:00"},
|
|
"dataset": {
|
|
"granularity": "None",
|
|
"aggregation": {"totalCost": {"name": "PreTaxCost", "function": "Sum"}},
|
|
"grouping": [{"type": "Dimension", "name": "ServiceName"}],
|
|
},
|
|
}
|
|
|
|
result = _azure_query_with_retry(client, scope, query)
|
|
|
|
columns = [c.name for c in result.columns]
|
|
cost_idx = columns.index("PreTaxCost")
|
|
currency_idx = columns.index("Currency") if "Currency" in columns else None
|
|
service_idx = columns.index("ServiceName") if "ServiceName" in columns else None
|
|
|
|
total = 0.0
|
|
by_service: dict[str, float] = {}
|
|
currency = "EUR"
|
|
for row in result.rows:
|
|
amount = float(row[cost_idx])
|
|
total += amount
|
|
if service_idx is not None:
|
|
service = row[service_idx]
|
|
by_service[service] = by_service.get(service, 0.0) + amount
|
|
if currency_idx is not None:
|
|
currency = row[currency_idx]
|
|
|
|
return {
|
|
"total": round(total, 2),
|
|
"currency": currency,
|
|
"by_service": {k: round(v, 2) for k, v in sorted(by_service.items(), key=lambda x: -x[1])},
|
|
}
|
|
|
|
|
|
def get_azure_resource_costs(tenant_id: str, client_id: str, client_secret: str,
|
|
subscription_id: str, start: str, end: str) -> dict:
|
|
"""Coût par ressource (VM, storage account...), groupé par service, sur
|
|
la période [start, end). Azure n'impose pas de limite de fenêtre pour ce
|
|
niveau de détail (contrairement à AWS) — on garde quand même 14 jours
|
|
en pratique pour rester cohérent entre les deux providers."""
|
|
from azure.identity import ClientSecretCredential
|
|
from azure.mgmt.costmanagement import CostManagementClient
|
|
from azure.core.exceptions import HttpResponseError
|
|
|
|
credential = ClientSecretCredential(tenant_id, client_id, client_secret)
|
|
client = CostManagementClient(credential)
|
|
scope = f"/subscriptions/{subscription_id}"
|
|
|
|
query = {
|
|
"type": "ActualCost",
|
|
"timeframe": "Custom",
|
|
"timePeriod": {"from": f"{start}T00:00:00+00:00", "to": f"{end}T00:00:00+00:00"},
|
|
"dataset": {
|
|
"granularity": "None",
|
|
"aggregation": {"totalCost": {"name": "PreTaxCost", "function": "Sum"}},
|
|
"grouping": [
|
|
{"type": "Dimension", "name": "ServiceName"},
|
|
{"type": "Dimension", "name": "ResourceId"},
|
|
],
|
|
},
|
|
}
|
|
|
|
try:
|
|
result = _azure_query_with_retry(client, scope, query)
|
|
except HttpResponseError as e:
|
|
raise ResourceLevelUnavailable(str(e)) from e
|
|
|
|
columns = [c.name for c in result.columns]
|
|
cost_idx = columns.index("PreTaxCost")
|
|
currency_idx = columns.index("Currency") if "Currency" in columns else None
|
|
service_idx = columns.index("ServiceName") if "ServiceName" in columns else None
|
|
resource_idx = columns.index("ResourceId") if "ResourceId" in columns else None
|
|
|
|
services: dict[str, dict[str, float]] = {}
|
|
currency = "EUR"
|
|
for row in result.rows:
|
|
amount = float(row[cost_idx])
|
|
if currency_idx is not None:
|
|
currency = row[currency_idx]
|
|
service = row[service_idx] if service_idx is not None else "(service inconnu)"
|
|
resource_id = row[resource_idx] if resource_idx is not None else ""
|
|
bucket = services.setdefault(service, {})
|
|
bucket[resource_id] = bucket.get(resource_id, 0.0) + amount
|
|
|
|
return {"currency": currency, "services": _finalize_resource_buckets(services)}
|
|
|
|
|
|
# ----------------------------------------------------------------------
|
|
# Conversion de devise
|
|
# ----------------------------------------------------------------------
|
|
|
|
_rate_cache: dict[tuple[str, str], float] = {}
|
|
|
|
|
|
def get_exchange_rate(from_currency: str, to_currency: str) -> float:
|
|
"""Taux de change du jour (source: BCE, via l'API gratuite frankfurter.app).
|
|
Mis en cache pour la durée du process."""
|
|
if from_currency == to_currency:
|
|
return 1.0
|
|
|
|
key = (from_currency, to_currency)
|
|
if key in _rate_cache:
|
|
return _rate_cache[key]
|
|
|
|
import requests
|
|
|
|
resp = requests.get(
|
|
"https://api.frankfurter.app/latest",
|
|
params={"from": from_currency, "to": to_currency},
|
|
timeout=10,
|
|
)
|
|
resp.raise_for_status()
|
|
rate = resp.json()["rates"][to_currency]
|
|
_rate_cache[key] = rate
|
|
return rate
|
|
|
|
|
|
def convert_result(result: dict, target_currency: str) -> dict:
|
|
"""Convertit un résultat (total + by_service) vers la devise cible.
|
|
Ajoute 'original_total' / 'original_currency' / 'exchange_rate' pour
|
|
traçabilité. Le montant est figé : il ne sera plus reconverti ensuite."""
|
|
original_currency = result["currency"]
|
|
|
|
if original_currency == target_currency:
|
|
result["original_total"] = result["total"]
|
|
result["original_currency"] = original_currency
|
|
result["exchange_rate"] = 1.0
|
|
return result
|
|
|
|
try:
|
|
rate = get_exchange_rate(original_currency, target_currency)
|
|
except Exception as e:
|
|
result["conversion_error"] = str(e)
|
|
return result
|
|
|
|
result["original_total"] = result["total"]
|
|
result["original_currency"] = original_currency
|
|
result["exchange_rate"] = rate
|
|
result["total"] = round(result["total"] * rate, 2)
|
|
result["currency"] = target_currency
|
|
if result.get("by_service"):
|
|
result["by_service"] = {k: round(v * rate, 2) for k, v in result["by_service"].items()}
|
|
return result
|
|
|
|
|
|
def convert_resource_costs(result: dict, target_currency: str) -> dict:
|
|
"""Équivalent de convert_result pour le résultat de
|
|
get_aws_resource_costs / get_azure_resource_costs (par ressource,
|
|
groupé par service)."""
|
|
original_currency = result["currency"]
|
|
if original_currency == target_currency:
|
|
return result
|
|
|
|
try:
|
|
rate = get_exchange_rate(original_currency, target_currency)
|
|
except Exception as e:
|
|
result["conversion_error"] = str(e)
|
|
return result
|
|
|
|
result["currency"] = target_currency
|
|
result["services"] = {
|
|
service: [{**r, "amount": round(r["amount"] * rate, 2)} for r in resources]
|
|
for service, resources in result["services"].items()
|
|
}
|
|
return result
|