first (big) draft
This commit is contained in:
@@ -0,0 +1,395 @@
|
||||
"""
|
||||
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
|
||||
Reference in New Issue
Block a user